From 34415fc817aceb976f3def56a351e034882f52cc Mon Sep 17 00:00:00 2001 From: Wyatt Verchere Date: Thu, 10 Sep 2026 12:27:08 -0700 Subject: [PATCH 01/12] feat: Adds the background worker --- .../integration-tests/tests/mars_async.rs | 318 ++++++++++++++++++ components/ads-client/src/ads_store.rs | 63 +++- components/ads-client/src/client.rs | 29 +- components/ads-client/src/client/error.rs | 34 +- components/ads-client/src/ffi.rs | 7 +- components/ads-client/src/lib.rs | 133 +++++++- components/ads-client/src/worker.rs | 63 ++++ components/ads-client/src/worker/command.rs | 190 +++++++++++ 8 files changed, 823 insertions(+), 14 deletions(-) create mode 100644 components/ads-client/integration-tests/tests/mars_async.rs create mode 100644 components/ads-client/src/worker.rs create mode 100644 components/ads-client/src/worker/command.rs diff --git a/components/ads-client/integration-tests/tests/mars_async.rs b/components/ads-client/integration-tests/tests/mars_async.rs new file mode 100644 index 00000000000..4ae8102d3ed --- /dev/null +++ b/components/ads-client/integration-tests/tests/mars_async.rs @@ -0,0 +1,318 @@ +/* This Source Code Form is subject to the terms of the Mozilla Public +* License, v. 2.0. If a copy of the MPL was not distributed with this +* file, You can obtain one at http://mozilla.org/MPL/2.0/. +*/ + +use ads_client::MozAdsIABContent; +use ads_client::MozAdsIABContentTaxonomy; +use ads_client::MozAdsPlacementRequestWithCount; +use ads_client::ads_store::PlacementId; +use ads_client::{ + MozAdsClient, MozAdsClientBuilder, MozAdsEnvironment, MozAdsPlacementRequest, + MozAdsRequestOptions, MozAdsTile, +}; +use std::sync::Arc; + +pub const TEST_TIMEOUT_DURATION: std::time::Duration = std::time::Duration::from_secs(10); + +fn init_backend() { + viaduct_hyper::viaduct_init_backend_hyper(); +} + +fn prod_client() -> ads_client::MozAdsClient { + Arc::new(MozAdsClientBuilder::new()) + .environment(MozAdsEnvironment::Prod) + .build() +} + +#[test] +#[ignore = "integration test: run manually with -- --ignored"] +fn test_contract_image_prod_async() { + init_backend(); + + // Prefetch + let placement_id = PlacementId::new("mock_billboard_1"); + let client = prod_client(); + let result = client.prefetch_ads( + vec![MozAdsPlacementRequest { + iab_content: None, + placement_id: placement_id.clone().into(), + }], + vec![], + vec![], + None, + ); + + assert!( + result.is_ok(), + "Image ad dispatch request failed: {:?}", + result.err() + ); + + // Ping (waits for queue to clear) + let ping = client.ping_background_worker(Some(TEST_TIMEOUT_DURATION)); + assert!(ping.is_ok(), "Ping failed: {:?}", ping.err()); + + // Query + let result = client.query_image_ads(placement_id); + assert!( + result.is_ok(), + "Querying for ads failed: {:?}", + result.err() + ); + let placements = result.unwrap(); + + assert!(placements.is_some()); +} + +#[test] +#[ignore = "integration test: run manually with -- --ignored"] +fn test_contract_image_with_categories_prod_async() { + init_backend(); + + // Prefetch + let placement_id = PlacementId::new("mock_billboard_1"); + let client = prod_client(); + let result = client.prefetch_ads( + vec![MozAdsPlacementRequest { + iab_content: Some(MozAdsIABContent { + category_ids: vec!["338".to_string()], + taxonomy: MozAdsIABContentTaxonomy::IAB3_0, + }), + placement_id: placement_id.clone().into(), + }], + vec![], + vec![], + Some(MozAdsRequestOptions { + flags: std::collections::HashMap::from([("contextual_placement".to_string(), true)]), + ..Default::default() + }), + ); + + assert!( + result.is_ok(), + "Image ad dispatch request failed: {:?}", + result.err() + ); + + // Ping (waits for queue to clear) + let ping = client.ping_background_worker(Some(TEST_TIMEOUT_DURATION)); + assert!(ping.is_ok(), "Ping failed: {:?}", ping.err()); + + // Query + let result = client.query_image_ads(placement_id); + assert!( + result.is_ok(), + "Querying for ads failed: {:?}", + result.err() + ); + + let placements = result.unwrap(); + assert!(placements.is_some()); + let ad = placements.unwrap(); + assert!(!ad.url.is_empty(), "destination url should be populated"); + assert!(!ad.image_url.is_empty(), "image url should be populated"); +} + +#[test] +#[ignore = "integration test: run manually with -- --ignored"] +fn test_contract_spoc_prod_async() { + init_backend(); + + // Prefetch + let placement_id = PlacementId::new("mock_spoc_1"); + let client = prod_client(); + let result = client.prefetch_ads( + vec![], + vec![MozAdsPlacementRequestWithCount { + count: 3, + iab_content: None, + placement_id: placement_id.clone().into(), + }], + vec![], + None, + ); + + assert!( + result.is_ok(), + "Spoc ad dispatch request failed: {:?}", + result.err() + ); + + // Ping (waits for queue to clear) + let ping = client.ping_background_worker(Some(TEST_TIMEOUT_DURATION)); + assert!(ping.is_ok(), "Ping failed: {:?}", ping.err()); + + // Query + let result = client.query_spoc_ads(placement_id); + assert!( + result.is_ok(), + "Querying for ads failed: {:?}", + result.err() + ); + let placements = result.unwrap(); + assert!(placements.is_some()); + assert!(placements.unwrap().len() == 3); +} + +#[test] +#[ignore = "integration test: run manually with -- --ignored"] +fn test_contract_tile_prod_async() { + init_backend(); + + // Prefetch + let placement_id = PlacementId::new("mock_tile_1"); + let client = prod_client(); + let result = client.prefetch_ads( + vec![], + vec![], + vec![MozAdsPlacementRequest { + iab_content: None, + placement_id: placement_id.clone().into(), + }], + None, + ); + + assert!( + result.is_ok(), + "Tile ad dispatch request failed: {:?}", + result.err() + ); + + // Ping (waits for queue to clear) + let ping = client.ping_background_worker(Some(TEST_TIMEOUT_DURATION)); + assert!(ping.is_ok(), "Ping failed: {:?}", ping.err()); + + // Query + let result = client.query_tile_ads(placement_id); + assert!( + result.is_ok(), + "Querying for ads failed: {:?}", + result.err() + ); + let placements = result.unwrap(); + + assert!(placements.is_some()); +} + +#[test] +#[ignore = "integration test: run manually with -- --ignored"] +fn test_contract_tile_ohttp_prod_async() { + init_backend(); + viaduct::ohttp::configure_ohttp_channel( + "ads-client".to_string(), + viaduct::ohttp::OhttpConfig { + relay_url: "https://mozilla-ohttp.fastly-edge.com/".to_string(), + gateway_host: "prod.ohttp-gateway.prod.webservices.mozgcp.net".to_string(), + }, + ) + .expect("OHTTP channel configuration should succeed"); + + // Prefetch + let placement_id = PlacementId::new("mock_tile_1"); + let client = prod_client(); + let result = client.prefetch_ads( + vec![], + vec![], + vec![MozAdsPlacementRequest { + iab_content: None, + placement_id: placement_id.clone().into(), + }], + Some(MozAdsRequestOptions { + ohttp: true, + ..Default::default() + }), + ); + + assert!( + result.is_ok(), + "Tile ad dispatch request failed: {:?}", + result.err() + ); + + // Ping (waits for queue to clear) + let ping = client.ping_background_worker(Some(TEST_TIMEOUT_DURATION)); + assert!(ping.is_ok(), "Ping failed: {:?}", ping.err()); + + // Query + let result = client.query_tile_ads(placement_id); + assert!( + result.is_ok(), + "Querying for ads failed: {:?}", + result.err() + ); + let placements = result.unwrap(); + + assert!( + placements.is_some(), + "OHTTP response should contain mock_tile_1" + ); +} + +#[test] +#[ignore = "integration test: run manually with -- --ignored"] +fn test_contract_multi_ad_type_prod_async() { + init_backend(); + + // Prefetch + let placement_image_id = PlacementId::new("mock_billboard_1"); + let placement_spoc_id = PlacementId::new("mock_spoc_1"); + let placement_tile_id = PlacementId::new("mock_tile_1"); + let client = prod_client(); + let result = client.prefetch_ads( + vec![MozAdsPlacementRequest { + iab_content: None, + placement_id: placement_image_id.clone().into(), + }], + vec![MozAdsPlacementRequestWithCount { + count: 4, + iab_content: None, + placement_id: placement_spoc_id.clone().into(), + }], + vec![MozAdsPlacementRequest { + iab_content: None, + placement_id: placement_tile_id.clone().into(), + }], + None, + ); + + assert!( + result.is_ok(), + "Image ad dispatch request failed: {:?}", + result.err() + ); + + // Ping (waits for queue to clear) + let ping = client.ping_background_worker(Some(TEST_TIMEOUT_DURATION)); + assert!(ping.is_ok(), "Ping failed: {:?}", ping.err()); + + // Query + let result = client.query_image_ads(placement_image_id); + assert!( + result.is_ok(), + "Querying for image ads failed: {:?}", + result.err() + ); + let placements = result.unwrap(); + assert!(placements.is_some()); + + let result = client.query_spoc_ads(placement_spoc_id); + assert!( + result.is_ok(), + "Querying for spoc ads failed: {:?}", + result.err() + ); + let placements = result.unwrap(); + assert!(placements.is_some()); + assert!(placements.unwrap().len() == 4); + + let result = client.query_tile_ads(placement_tile_id); + assert!( + result.is_ok(), + "Querying for ads failed: {:?}", + result.err() + ); + let placements = result.unwrap(); + + assert!(placements.is_some()); +} \ No newline at end of file diff --git a/components/ads-client/src/ads_store.rs b/components/ads-client/src/ads_store.rs index eb8c377d94d..00b7861adf6 100644 --- a/components/ads-client/src/ads_store.rs +++ b/components/ads-client/src/ads_store.rs @@ -7,9 +7,12 @@ use serde::{Deserialize, Serialize}; use crate::{ ads_store::{builder::AdsStoreBuilder, store::AdsStoreHolder}, common::bytesize::ByteSize, - mars::ad_response::{AdImage, AdSpoc, AdTile}, + mars::{ + ad_response::{AdImage, AdSpoc, AdTile}, + error::FetchAdsError, + }, }; -use std::path::Path; +use std::{collections::HashMap, path::Path}; /// Identification of placement sent and returned from MARS (eg: `mock_spoc_1`) #[derive(Debug, Hash, PartialEq, Eq, Clone)] @@ -19,9 +22,6 @@ impl PlacementId { pub fn new(s: &str) -> PlacementId { PlacementId(s.to_string()) } - pub fn into_inner(self) -> String { - self.0 - } } impl AsRef for PlacementId { @@ -30,13 +30,52 @@ impl AsRef for PlacementId { } } +impl From for PlacementId { + fn from(value: String) -> Self { + PlacementId(value) + } +} + +impl From for String { + fn from(value: PlacementId) -> Self { + value.0 + } +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] pub enum StorableAd { Image(AdImage), - Spoc(AdSpoc), + Spoc(Vec), Tile(AdTile), } +impl StorableAd { + // TODO: Should these be result? + pub fn into_image(self) -> Option { + if let StorableAd::Image(image) = self { + Some(image) + } else { + None + } + } + + pub fn into_spocs(self) -> Option> { + if let StorableAd::Spoc(spocs) = self { + Some(spocs) + } else { + None + } + } + + pub fn into_tile(self) -> Option { + if let StorableAd::Tile(tile) = self { + Some(tile) + } else { + None + } + } +} + pub struct AdsStore { holder: AdsStoreHolder, #[allow(dead_code)] @@ -61,6 +100,18 @@ impl AdsStore { self.holder.invalidate_ad_by_id(placement_id)?; Ok(()) } + + pub fn lookup(&self, placement_id: &PlacementId) -> Result, FetchAdsError> { + self.holder.lookup(placement_id) + } + + pub fn store_ads(&self, ads: HashMap) -> Result<(), FetchAdsError> { + // TODO: The query can be rewritten to have multiple ads inserted in one + for (placement_id, ad) in ads { + self.holder.store_ad(&placement_id, ad)?; + } + Ok(()) + } } #[cfg(test)] diff --git a/components/ads-client/src/client.rs b/components/ads-client/src/client.rs index 99290582119..739892f6fc5 100644 --- a/components/ads-client/src/client.rs +++ b/components/ads-client/src/client.rs @@ -7,12 +7,12 @@ use std::collections::HashMap; use std::sync::Arc; use std::time::Duration; -use crate::ads_store::AdsStore; +use crate::ads_store::{AdsStore, PlacementId, StorableAd}; use crate::common::bytesize::ByteSize; use crate::http_cache::{CachePolicy, HttpCache}; use crate::mars::ad_request::{AdPlacementRequest, AdRequestFlags}; use crate::mars::ad_response::{AdImage, AdResponse, AdResponseValue, AdSpoc, AdTile}; -use crate::mars::error::{RecordClickError, RecordImpressionError, ReportAdError}; +use crate::mars::error::{FetchAdsError, RecordClickError, RecordImpressionError, ReportAdError}; use crate::mars::{MARSClient, ReportReason}; use crate::shutdown::{AdsStoreShutdown, ShutdownReferences}; use crate::telemetry::Telemetry; @@ -116,6 +116,31 @@ where self.client.clear_cache() } + pub fn cache_ads( + &mut self, + ads: HashMap, + ) -> Result<(), FetchAdsError> { + let ads_store = self.ads_store.lock(); + if let Some(ads_store) = ads_store.as_ref() { + ads_store.store_ads(ads)?; + } + // TODO: Should we error if no ads store? + Ok(()) + } + + pub fn get_cached_ad( + &self, + placement_id: &PlacementId, + ) -> Result, FetchAdsError> { + let ads_store = self.ads_store.lock(); + if let Some(ads_store) = ads_store.as_ref() { + Ok(ads_store.lookup(placement_id)?) + } else { + // TODO: Should we error if no ads store? + Ok(None) + } + } + pub fn get_context_id(&self) -> context_id::ApiResult { self.context_id_provider.context_id() } diff --git a/components/ads-client/src/client/error.rs b/components/ads-client/src/client/error.rs index 2542939493f..913a36792c9 100644 --- a/components/ads-client/src/client/error.rs +++ b/components/ads-client/src/client/error.rs @@ -3,10 +3,18 @@ * file, You can obtain one at http://mozilla.org/MPL/2.0/. */ -use crate::mars::error::{FetchAdsError, RecordClickError, RecordImpressionError, ReportAdError}; +use std::sync::mpsc::{RecvTimeoutError, TrySendError}; + +use crate::{ + mars::error::{FetchAdsError, RecordClickError, RecordImpressionError, ReportAdError}, + worker::command, +}; #[derive(Debug, thiserror::Error)] pub enum ComponentError { + #[error("Error requesting ads from worker: {0}")] + BackgroundWorker(#[from] BackgroundWorkerError), + #[error("Error recording a click for a placement: {0}")] RecordClick(#[from] RecordClickError), @@ -28,3 +36,27 @@ pub enum RequestAdsError { #[error("Error requesting ads from MARS: {0}")] FetchAds(#[from] FetchAdsError), } + +#[derive(Debug, thiserror::Error)] +pub enum BackgroundWorkerError { + #[error("Error sending pong back from background worker")] + PongFailure(Box>), + + #[error("Error requesting new ads from the background worker: worker closed")] + WorkerClosed, + + #[error("Error requesting new ads from the background worker: worker full")] + WorkerFull, + + #[error("Worker timed out waiting for response: {0}")] + WorkerTimedOut(#[from] RecvTimeoutError), +} + +impl From> for BackgroundWorkerError { + fn from(value: TrySendError) -> Self { + match value { + TrySendError::Disconnected(_) => BackgroundWorkerError::WorkerClosed, + TrySendError::Full(_) => BackgroundWorkerError::WorkerFull, + } + } +} diff --git a/components/ads-client/src/ffi.rs b/components/ads-client/src/ffi.rs index 4cc40fa8a82..c17584a1238 100644 --- a/components/ads-client/src/ffi.rs +++ b/components/ads-client/src/ffi.rs @@ -23,7 +23,7 @@ use crate::mars::ad_response::{ use crate::mars::Environment; use crate::mars::ReportReason; use crate::AdsClientUrl; -use crate::MozAdsClient; +use crate::{worker, MozAdsClient}; use parking_lot::Mutex; use std::collections::HashMap; @@ -147,9 +147,12 @@ impl MozAdsClientBuilder { }; let client = AdsClient::new(client_config); let shutdown_references = client.shutdown_references(); + let inner = Arc::new(Mutex::new(client)); + let worker = worker::AdsClientWorkerWrapper::new(inner.clone()); MozAdsClient { - inner: Mutex::new(client), + inner, shutdown_references, + worker, } } diff --git a/components/ads-client/src/lib.rs b/components/ads-client/src/lib.rs index a599362b5c1..63832b8bfd4 100644 --- a/components/ads-client/src/lib.rs +++ b/components/ads-client/src/lib.rs @@ -3,7 +3,11 @@ * file, You can obtain one at http://mozilla.org/MPL/2.0/. */ -use std::collections::HashMap; +use std::{ + collections::HashMap, + sync::{mpsc, Arc}, + time::Duration, +}; use client::error::ComponentError; use error_support::handle_error; @@ -22,10 +26,17 @@ pub mod http_cache; mod mars; pub mod shutdown; pub mod telemetry; +pub mod worker; pub use ffi::*; -use crate::{ffi::telemetry::MozAdsTelemetryWrapper, shutdown::ShutdownReferences}; +use crate::{ + ads_store::{PlacementId, StorableAd}, + client::error::{BackgroundWorkerError, RequestAdsError}, + ffi::telemetry::MozAdsTelemetryWrapper, + shutdown::ShutdownReferences, + worker::{command::DispatchCommand, AdsClientWorkerWrapper}, +}; #[cfg(test)] mod test_utils; @@ -37,11 +48,14 @@ uniffi::custom_type!(AdsClientUrl, String, { try_lift: |val| Ok(AdsClientUrl::parse(&val)?), lower: |obj| obj.as_str().to_string(), }); +uniffi::custom_type!(PlacementId, String); +pub type MozAdsClientInner = Arc>>; #[derive(uniffi::Object)] pub struct MozAdsClient { - inner: Mutex>, + inner: MozAdsClientInner, shutdown_references: ShutdownReferences, + worker: AdsClientWorkerWrapper, } #[uniffi::export] @@ -179,4 +193,117 @@ impl MozAdsClient { .map_err(ComponentError::RequestAds)?; Ok(response.into_iter().map(|(k, v)| (k, v.into())).collect()) } + + // TODO: Can we make this one request? + #[handle_error(ComponentError)] + #[uniffi::method(default(image_ad_requests = [], spoc_ad_requests = [], tile_ad_requests = [], options = None))] + pub fn prefetch_ads( + &self, + image_ad_requests: Vec, + spoc_ad_requests: Vec, + tile_ad_requests: Vec, + options: Option, + ) -> AdsClientApiResult<()> { + let options = options.unwrap_or_default(); + let flags = AdRequestFlags::from(&options); + let ohttp = options.ohttp; + let blocks = options.blocks.clone(); + let cache_policy: CachePolicy = options.into(); + + // Dispatch image requests + if !image_ad_requests.is_empty() { + self.worker.dispatch(DispatchCommand::RequestImageAds { + image_ad_requests, + ohttp, + cache_policy, + flags: flags.clone(), + blocks: blocks.clone(), + })?; + } + // Dispatch spoc requests + if !spoc_ad_requests.is_empty() { + self.worker.dispatch(DispatchCommand::RequestSpocAds { + spoc_ad_requests, + ohttp, + cache_policy, + flags: flags.clone(), + blocks: blocks.clone(), + })?; + } + + // Dispatch tiles requests + if !tile_ad_requests.is_empty() { + self.worker.dispatch(DispatchCommand::RequestTileAds { + tile_ad_requests, + ohttp, + cache_policy, + flags: flags.clone(), + blocks: blocks.clone(), + })?; + } + + Ok(()) + } + + #[handle_error(ComponentError)] + #[uniffi::method()] + pub fn query_image_ads( + &self, + placement_id: PlacementId, + ) -> AdsClientApiResult> { + let inner = self.inner.lock(); + let image_ad: Option = inner + .get_cached_ad(&placement_id) + .map_err(RequestAdsError::from)?; + Ok(image_ad + .and_then(|ad| ad.into_image()) + .map(|ad| ad.clone().into())) + } + + #[handle_error(ComponentError)] + #[uniffi::method()] + pub fn query_spoc_ads( + &self, + placement_id: PlacementId, + ) -> AdsClientApiResult>> { + let inner = self.inner.lock(); + let spoc_ads: Option = inner + .get_cached_ad(&placement_id) + .map_err(RequestAdsError::from)?; + Ok(spoc_ads + .and_then(|ad| ad.into_spocs()) + .map(|ad| ad.clone().into_iter().map(|x| x.into()).collect())) + } + + #[handle_error(ComponentError)] + #[uniffi::method()] + pub fn query_tile_ads( + &self, + placement_id: PlacementId, + ) -> AdsClientApiResult> { + let inner = self.inner.lock(); + let tile_ad: Option = inner + .get_cached_ad(&placement_id) + .map_err(RequestAdsError::from)?; + Ok(tile_ad + .and_then(|ad| ad.into_tile()) + .map(|ad| ad.clone().into())) + } + + // Pings the background worker and waits for a response back, for use in tests. + // Because the background worker is synchronous, this returns if the worker is empty, + // making it useful for integration tests to wait until all tasks have completed. + #[handle_error(ComponentError)] + pub fn ping_background_worker(&self, timeout: Option) -> AdsClientApiResult<()> { + let (tx, rx) = mpsc::sync_channel(0); + self.worker.dispatch(DispatchCommand::Ping(tx))?; + + if let Some(timeout) = timeout { + rx.recv_timeout(timeout) + .map_err(BackgroundWorkerError::from)?; + } else { + rx.recv().map_err(|_| BackgroundWorkerError::WorkerClosed)?; + } + Ok(()) + } } diff --git a/components/ads-client/src/worker.rs b/components/ads-client/src/worker.rs new file mode 100644 index 00000000000..0538d7e99b9 --- /dev/null +++ b/components/ads-client/src/worker.rs @@ -0,0 +1,63 @@ +use crate::{ + client::error::{BackgroundWorkerError, ComponentError}, + worker::command::DispatchCommand, + MozAdsClientInner, +}; +use std::{ + sync::mpsc::{self, Receiver, SyncSender}, + thread::JoinHandle, +}; + +pub mod command; + +pub const ADS_CLIENT_WORKER_CHANNEL_BUFFER_SIZE: usize = 1000; +pub const ADS_CLIENT_WORKER_THREAD_NAME: &str = "ads-client.worker"; + +pub struct AdsClientWorkerWrapper { + _worker_thread: Option>, + worker_dispatch: Option>, +} + +impl AdsClientWorkerWrapper { + pub fn new(inner: MozAdsClientInner) -> AdsClientWorkerWrapper { + let (worker_dispatch, worker_thread) = Option::unzip(build_worker_thread(inner.clone())); + AdsClientWorkerWrapper { + _worker_thread: worker_thread, + worker_dispatch, + } + } + + pub fn dispatch(&self, command: DispatchCommand) -> Result<(), ComponentError> { + if let Some(worker_dispatch) = &self.worker_dispatch { + worker_dispatch + .try_send(command) + .map_err(BackgroundWorkerError::from)?; + + Ok(()) + } else { + Err(BackgroundWorkerError::WorkerClosed.into()) + } + } +} + +// Spawn worker thread from a reference to the client, returning a synchronous channel transmitter to the thread, and its JoinHandle. +// Returns None if thread fails to build. +pub fn build_worker_thread( + inner_client: MozAdsClientInner, +) -> Option<(SyncSender, JoinHandle<()>)> { + let (tx, rx) = mpsc::sync_channel(ADS_CLIENT_WORKER_CHANNEL_BUFFER_SIZE); + let worker_thread_handle = std::thread::Builder::new() + .name(ADS_CLIENT_WORKER_THREAD_NAME.to_string()) + .spawn(move || crate::worker::worker(inner_client, rx)).inspect_err(|err| { + error_support::error!("Failed to create ads-client worker thread `{ADS_CLIENT_WORKER_THREAD_NAME}` with: {err}") + }).ok()?; + Some((tx, worker_thread_handle)) +} + +fn worker(inner_client: MozAdsClientInner, rx: Receiver) { + // Synchronously run tasks in the order they are passed in this separate channel. + while let Ok(command) = rx.recv() { + // Error is naturally logged through `handle_error` conversion macro. + let _ = command.run_command(&inner_client); + } +} diff --git a/components/ads-client/src/worker/command.rs b/components/ads-client/src/worker/command.rs new file mode 100644 index 00000000000..a1bb1ebba4a --- /dev/null +++ b/components/ads-client/src/worker/command.rs @@ -0,0 +1,190 @@ +use std::{collections::HashMap, sync::mpsc::SyncSender}; + +use error_support::handle_error; +use url::Url; + +use crate::{ + ads_store::StorableAd, + client::error::{BackgroundWorkerError, ComponentError, RequestAdsError}, + http_cache::CachePolicy, + mars::{ad_request::AdPlacementRequest, ReportReason}, + AdsClientApiResult, MozAdsClientInner, MozAdsPlacementRequest, MozAdsPlacementRequestWithCount, +}; + +// Command dispatch enum for passing different instructions to the background worker thread. +// `RequestImageAds`, `RequestSpocAds`, `RequestTileAds` are prefetch mechanisms that query and load data into the local cache. +// `RecordClick`, `RecordImpression`, `ReportAd` are fire and forget mechanisms that do not load data. +// `Ping` is a synchronous command for internal use that triggers its inner channel when command resolves (eg: when the queue is empty). +pub enum DispatchCommand { + RequestImageAds { + image_ad_requests: Vec, + cache_policy: CachePolicy, + ohttp: bool, + flags: HashMap, + blocks: Vec, + }, + RequestSpocAds { + spoc_ad_requests: Vec, + cache_policy: CachePolicy, + ohttp: bool, + flags: HashMap, + blocks: Vec, + }, + RequestTileAds { + tile_ad_requests: Vec, + cache_policy: CachePolicy, + ohttp: bool, + flags: HashMap, + blocks: Vec, + }, + RecordClick { + url: Url, + ohttp: bool, + }, + RecordImpression { + url: Url, + ohttp: bool, + }, + ReportAd { + url: Url, + reason: ReportReason, + ohttp: bool, + }, + Ping(SyncSender<()>), +} + +impl DispatchCommand { + // Runs a dispatched command synchronously in it's thread. + // The dispatched command calls the corresponding `AdsClient` synchronous method, meaning that behavior between the two is shared. + // This includes telemetry calls, meaning that for a successful `RecordClick`, all of the following will get logged: + // - CommandDispatchedOperationEvent::RecordClick (on dispatch) + // - ClientOperationEvent::RecordClick (on `AdsClient` method success) + // - CommandProcessedOperationEvent::RecordClick (on process) + #[handle_error(ComponentError)] + pub fn run_command(self, ads_client_inner: &MozAdsClientInner) -> AdsClientApiResult<()> { + match self { + DispatchCommand::RequestImageAds { + image_ad_requests, + cache_policy, + flags, + ohttp, + blocks, + } => { + let mut inner = ads_client_inner.lock(); + if !image_ad_requests.is_empty() { + let image_ad_requests: Vec = + image_ad_requests.iter().map(|r| r.into()).collect(); + let image_response = inner + .request_image_ads( + image_ad_requests, + flags, + Some(cache_policy), + ohttp, + blocks, + ) + .map_err(ComponentError::RequestAds)?; + // TODO: Bulk insert + inner + .cache_ads( + image_response + .into_iter() + .map(|(k, v)| (k.into(), StorableAd::Image(v))) + .collect(), + ) + .map_err(RequestAdsError::from)?; + } + Ok(()) + } + // TODO: Can we modify these to be one call? + DispatchCommand::RequestSpocAds { + spoc_ad_requests, + cache_policy, + flags, + ohttp, + blocks, + } => { + let mut inner = ads_client_inner.lock(); + if !spoc_ad_requests.is_empty() { + let spoc_ad_requests: Vec = + spoc_ad_requests.iter().map(|r| r.into()).collect(); + let spoc_response = inner + .request_spoc_ads( + spoc_ad_requests, + flags, + Some(cache_policy), + ohttp, + blocks, + ) + .map_err(ComponentError::RequestAds)?; + inner + .cache_ads( + spoc_response + .into_iter() + .map(|(k, v)| (k.into(), StorableAd::Spoc(v))) + .collect(), + ) + .map_err(RequestAdsError::from)?; + } + Ok(()) + } + DispatchCommand::RequestTileAds { + tile_ad_requests, + cache_policy, + flags, + ohttp, + blocks, + } => { + let mut inner = ads_client_inner.lock(); + if !tile_ad_requests.is_empty() { + let tile_ad_requests: Vec = + tile_ad_requests.iter().map(|r| r.into()).collect(); + let tile_response = inner + .request_tile_ads( + tile_ad_requests, + flags, + Some(cache_policy), + ohttp, + blocks, + ) + .map_err(ComponentError::RequestAds)?; + inner + .cache_ads( + tile_response + .into_iter() + .map(|(k, v)| (k.into(), StorableAd::Tile(v))) + .collect(), + ) + .map_err(RequestAdsError::from)?; + } + Ok(()) + } + DispatchCommand::RecordClick { url, ohttp } => { + let inner = ads_client_inner.lock(); + inner + .record_click(url, ohttp) + .map_err(ComponentError::RecordClick)?; + Ok(()) + } + DispatchCommand::RecordImpression { url, ohttp } => { + let inner = ads_client_inner.lock(); + inner + .record_impression(url, ohttp) + .map_err(ComponentError::RecordImpression)?; + Ok(()) + } + DispatchCommand::ReportAd { url, ohttp, reason } => { + let inner = ads_client_inner.lock(); + inner + .report_ad(url, reason, ohttp) + .map_err(ComponentError::ReportAd)?; + Ok(()) + } + DispatchCommand::Ping(sender) => { + sender + .try_send(()) + .map_err(|err| BackgroundWorkerError::PongFailure(Box::new(err)))?; + Ok(()) + } + } + } +} From 791f0a6d2c5edb551fb1a7d9d2e68e95f27b5888 Mon Sep 17 00:00:00 2001 From: Wyatt Verchere Date: Thu, 10 Sep 2026 13:11:47 -0700 Subject: [PATCH 02/12] fix: removes the FFI layer --- .../integration-tests/tests/mars_async.rs | 318 ------------------ components/ads-client/src/lib.rs | 100 +----- 2 files changed, 1 insertion(+), 417 deletions(-) delete mode 100644 components/ads-client/integration-tests/tests/mars_async.rs diff --git a/components/ads-client/integration-tests/tests/mars_async.rs b/components/ads-client/integration-tests/tests/mars_async.rs deleted file mode 100644 index 4ae8102d3ed..00000000000 --- a/components/ads-client/integration-tests/tests/mars_async.rs +++ /dev/null @@ -1,318 +0,0 @@ -/* This Source Code Form is subject to the terms of the Mozilla Public -* License, v. 2.0. If a copy of the MPL was not distributed with this -* file, You can obtain one at http://mozilla.org/MPL/2.0/. -*/ - -use ads_client::MozAdsIABContent; -use ads_client::MozAdsIABContentTaxonomy; -use ads_client::MozAdsPlacementRequestWithCount; -use ads_client::ads_store::PlacementId; -use ads_client::{ - MozAdsClient, MozAdsClientBuilder, MozAdsEnvironment, MozAdsPlacementRequest, - MozAdsRequestOptions, MozAdsTile, -}; -use std::sync::Arc; - -pub const TEST_TIMEOUT_DURATION: std::time::Duration = std::time::Duration::from_secs(10); - -fn init_backend() { - viaduct_hyper::viaduct_init_backend_hyper(); -} - -fn prod_client() -> ads_client::MozAdsClient { - Arc::new(MozAdsClientBuilder::new()) - .environment(MozAdsEnvironment::Prod) - .build() -} - -#[test] -#[ignore = "integration test: run manually with -- --ignored"] -fn test_contract_image_prod_async() { - init_backend(); - - // Prefetch - let placement_id = PlacementId::new("mock_billboard_1"); - let client = prod_client(); - let result = client.prefetch_ads( - vec![MozAdsPlacementRequest { - iab_content: None, - placement_id: placement_id.clone().into(), - }], - vec![], - vec![], - None, - ); - - assert!( - result.is_ok(), - "Image ad dispatch request failed: {:?}", - result.err() - ); - - // Ping (waits for queue to clear) - let ping = client.ping_background_worker(Some(TEST_TIMEOUT_DURATION)); - assert!(ping.is_ok(), "Ping failed: {:?}", ping.err()); - - // Query - let result = client.query_image_ads(placement_id); - assert!( - result.is_ok(), - "Querying for ads failed: {:?}", - result.err() - ); - let placements = result.unwrap(); - - assert!(placements.is_some()); -} - -#[test] -#[ignore = "integration test: run manually with -- --ignored"] -fn test_contract_image_with_categories_prod_async() { - init_backend(); - - // Prefetch - let placement_id = PlacementId::new("mock_billboard_1"); - let client = prod_client(); - let result = client.prefetch_ads( - vec![MozAdsPlacementRequest { - iab_content: Some(MozAdsIABContent { - category_ids: vec!["338".to_string()], - taxonomy: MozAdsIABContentTaxonomy::IAB3_0, - }), - placement_id: placement_id.clone().into(), - }], - vec![], - vec![], - Some(MozAdsRequestOptions { - flags: std::collections::HashMap::from([("contextual_placement".to_string(), true)]), - ..Default::default() - }), - ); - - assert!( - result.is_ok(), - "Image ad dispatch request failed: {:?}", - result.err() - ); - - // Ping (waits for queue to clear) - let ping = client.ping_background_worker(Some(TEST_TIMEOUT_DURATION)); - assert!(ping.is_ok(), "Ping failed: {:?}", ping.err()); - - // Query - let result = client.query_image_ads(placement_id); - assert!( - result.is_ok(), - "Querying for ads failed: {:?}", - result.err() - ); - - let placements = result.unwrap(); - assert!(placements.is_some()); - let ad = placements.unwrap(); - assert!(!ad.url.is_empty(), "destination url should be populated"); - assert!(!ad.image_url.is_empty(), "image url should be populated"); -} - -#[test] -#[ignore = "integration test: run manually with -- --ignored"] -fn test_contract_spoc_prod_async() { - init_backend(); - - // Prefetch - let placement_id = PlacementId::new("mock_spoc_1"); - let client = prod_client(); - let result = client.prefetch_ads( - vec![], - vec![MozAdsPlacementRequestWithCount { - count: 3, - iab_content: None, - placement_id: placement_id.clone().into(), - }], - vec![], - None, - ); - - assert!( - result.is_ok(), - "Spoc ad dispatch request failed: {:?}", - result.err() - ); - - // Ping (waits for queue to clear) - let ping = client.ping_background_worker(Some(TEST_TIMEOUT_DURATION)); - assert!(ping.is_ok(), "Ping failed: {:?}", ping.err()); - - // Query - let result = client.query_spoc_ads(placement_id); - assert!( - result.is_ok(), - "Querying for ads failed: {:?}", - result.err() - ); - let placements = result.unwrap(); - assert!(placements.is_some()); - assert!(placements.unwrap().len() == 3); -} - -#[test] -#[ignore = "integration test: run manually with -- --ignored"] -fn test_contract_tile_prod_async() { - init_backend(); - - // Prefetch - let placement_id = PlacementId::new("mock_tile_1"); - let client = prod_client(); - let result = client.prefetch_ads( - vec![], - vec![], - vec![MozAdsPlacementRequest { - iab_content: None, - placement_id: placement_id.clone().into(), - }], - None, - ); - - assert!( - result.is_ok(), - "Tile ad dispatch request failed: {:?}", - result.err() - ); - - // Ping (waits for queue to clear) - let ping = client.ping_background_worker(Some(TEST_TIMEOUT_DURATION)); - assert!(ping.is_ok(), "Ping failed: {:?}", ping.err()); - - // Query - let result = client.query_tile_ads(placement_id); - assert!( - result.is_ok(), - "Querying for ads failed: {:?}", - result.err() - ); - let placements = result.unwrap(); - - assert!(placements.is_some()); -} - -#[test] -#[ignore = "integration test: run manually with -- --ignored"] -fn test_contract_tile_ohttp_prod_async() { - init_backend(); - viaduct::ohttp::configure_ohttp_channel( - "ads-client".to_string(), - viaduct::ohttp::OhttpConfig { - relay_url: "https://mozilla-ohttp.fastly-edge.com/".to_string(), - gateway_host: "prod.ohttp-gateway.prod.webservices.mozgcp.net".to_string(), - }, - ) - .expect("OHTTP channel configuration should succeed"); - - // Prefetch - let placement_id = PlacementId::new("mock_tile_1"); - let client = prod_client(); - let result = client.prefetch_ads( - vec![], - vec![], - vec![MozAdsPlacementRequest { - iab_content: None, - placement_id: placement_id.clone().into(), - }], - Some(MozAdsRequestOptions { - ohttp: true, - ..Default::default() - }), - ); - - assert!( - result.is_ok(), - "Tile ad dispatch request failed: {:?}", - result.err() - ); - - // Ping (waits for queue to clear) - let ping = client.ping_background_worker(Some(TEST_TIMEOUT_DURATION)); - assert!(ping.is_ok(), "Ping failed: {:?}", ping.err()); - - // Query - let result = client.query_tile_ads(placement_id); - assert!( - result.is_ok(), - "Querying for ads failed: {:?}", - result.err() - ); - let placements = result.unwrap(); - - assert!( - placements.is_some(), - "OHTTP response should contain mock_tile_1" - ); -} - -#[test] -#[ignore = "integration test: run manually with -- --ignored"] -fn test_contract_multi_ad_type_prod_async() { - init_backend(); - - // Prefetch - let placement_image_id = PlacementId::new("mock_billboard_1"); - let placement_spoc_id = PlacementId::new("mock_spoc_1"); - let placement_tile_id = PlacementId::new("mock_tile_1"); - let client = prod_client(); - let result = client.prefetch_ads( - vec![MozAdsPlacementRequest { - iab_content: None, - placement_id: placement_image_id.clone().into(), - }], - vec![MozAdsPlacementRequestWithCount { - count: 4, - iab_content: None, - placement_id: placement_spoc_id.clone().into(), - }], - vec![MozAdsPlacementRequest { - iab_content: None, - placement_id: placement_tile_id.clone().into(), - }], - None, - ); - - assert!( - result.is_ok(), - "Image ad dispatch request failed: {:?}", - result.err() - ); - - // Ping (waits for queue to clear) - let ping = client.ping_background_worker(Some(TEST_TIMEOUT_DURATION)); - assert!(ping.is_ok(), "Ping failed: {:?}", ping.err()); - - // Query - let result = client.query_image_ads(placement_image_id); - assert!( - result.is_ok(), - "Querying for image ads failed: {:?}", - result.err() - ); - let placements = result.unwrap(); - assert!(placements.is_some()); - - let result = client.query_spoc_ads(placement_spoc_id); - assert!( - result.is_ok(), - "Querying for spoc ads failed: {:?}", - result.err() - ); - let placements = result.unwrap(); - assert!(placements.is_some()); - assert!(placements.unwrap().len() == 4); - - let result = client.query_tile_ads(placement_tile_id); - assert!( - result.is_ok(), - "Querying for ads failed: {:?}", - result.err() - ); - let placements = result.unwrap(); - - assert!(placements.is_some()); -} \ No newline at end of file diff --git a/components/ads-client/src/lib.rs b/components/ads-client/src/lib.rs index 63832b8bfd4..f5a0dc71078 100644 --- a/components/ads-client/src/lib.rs +++ b/components/ads-client/src/lib.rs @@ -31,8 +31,7 @@ pub mod worker; pub use ffi::*; use crate::{ - ads_store::{PlacementId, StorableAd}, - client::error::{BackgroundWorkerError, RequestAdsError}, + client::error::BackgroundWorkerError, ffi::telemetry::MozAdsTelemetryWrapper, shutdown::ShutdownReferences, worker::{command::DispatchCommand, AdsClientWorkerWrapper}, @@ -48,7 +47,6 @@ uniffi::custom_type!(AdsClientUrl, String, { try_lift: |val| Ok(AdsClientUrl::parse(&val)?), lower: |obj| obj.as_str().to_string(), }); -uniffi::custom_type!(PlacementId, String); pub type MozAdsClientInner = Arc>>; #[derive(uniffi::Object)] @@ -194,102 +192,6 @@ impl MozAdsClient { Ok(response.into_iter().map(|(k, v)| (k, v.into())).collect()) } - // TODO: Can we make this one request? - #[handle_error(ComponentError)] - #[uniffi::method(default(image_ad_requests = [], spoc_ad_requests = [], tile_ad_requests = [], options = None))] - pub fn prefetch_ads( - &self, - image_ad_requests: Vec, - spoc_ad_requests: Vec, - tile_ad_requests: Vec, - options: Option, - ) -> AdsClientApiResult<()> { - let options = options.unwrap_or_default(); - let flags = AdRequestFlags::from(&options); - let ohttp = options.ohttp; - let blocks = options.blocks.clone(); - let cache_policy: CachePolicy = options.into(); - - // Dispatch image requests - if !image_ad_requests.is_empty() { - self.worker.dispatch(DispatchCommand::RequestImageAds { - image_ad_requests, - ohttp, - cache_policy, - flags: flags.clone(), - blocks: blocks.clone(), - })?; - } - // Dispatch spoc requests - if !spoc_ad_requests.is_empty() { - self.worker.dispatch(DispatchCommand::RequestSpocAds { - spoc_ad_requests, - ohttp, - cache_policy, - flags: flags.clone(), - blocks: blocks.clone(), - })?; - } - - // Dispatch tiles requests - if !tile_ad_requests.is_empty() { - self.worker.dispatch(DispatchCommand::RequestTileAds { - tile_ad_requests, - ohttp, - cache_policy, - flags: flags.clone(), - blocks: blocks.clone(), - })?; - } - - Ok(()) - } - - #[handle_error(ComponentError)] - #[uniffi::method()] - pub fn query_image_ads( - &self, - placement_id: PlacementId, - ) -> AdsClientApiResult> { - let inner = self.inner.lock(); - let image_ad: Option = inner - .get_cached_ad(&placement_id) - .map_err(RequestAdsError::from)?; - Ok(image_ad - .and_then(|ad| ad.into_image()) - .map(|ad| ad.clone().into())) - } - - #[handle_error(ComponentError)] - #[uniffi::method()] - pub fn query_spoc_ads( - &self, - placement_id: PlacementId, - ) -> AdsClientApiResult>> { - let inner = self.inner.lock(); - let spoc_ads: Option = inner - .get_cached_ad(&placement_id) - .map_err(RequestAdsError::from)?; - Ok(spoc_ads - .and_then(|ad| ad.into_spocs()) - .map(|ad| ad.clone().into_iter().map(|x| x.into()).collect())) - } - - #[handle_error(ComponentError)] - #[uniffi::method()] - pub fn query_tile_ads( - &self, - placement_id: PlacementId, - ) -> AdsClientApiResult> { - let inner = self.inner.lock(); - let tile_ad: Option = inner - .get_cached_ad(&placement_id) - .map_err(RequestAdsError::from)?; - Ok(tile_ad - .and_then(|ad| ad.into_tile()) - .map(|ad| ad.clone().into())) - } - // Pings the background worker and waits for a response back, for use in tests. // Because the background worker is synchronous, this returns if the worker is empty, // making it useful for integration tests to wait until all tasks have completed. From ef7a00ad5725f4ddb5afd5deacb2db77732499ba Mon Sep 17 00:00:00 2001 From: Wyatt Verchere Date: Thu, 10 Sep 2026 15:32:09 -0700 Subject: [PATCH 03/12] feat: fixes, updates --- components/ads-client/src/ads_store.rs | 1 - components/ads-client/src/ffi.rs | 7 +++- components/ads-client/src/worker.rs | 7 ++++ components/ads-client/src/worker/command.rs | 38 +-------------------- 4 files changed, 14 insertions(+), 39 deletions(-) diff --git a/components/ads-client/src/ads_store.rs b/components/ads-client/src/ads_store.rs index 00b7861adf6..bf72b36982d 100644 --- a/components/ads-client/src/ads_store.rs +++ b/components/ads-client/src/ads_store.rs @@ -106,7 +106,6 @@ impl AdsStore { } pub fn store_ads(&self, ads: HashMap) -> Result<(), FetchAdsError> { - // TODO: The query can be rewritten to have multiple ads inserted in one for (placement_id, ad) in ads { self.holder.store_ad(&placement_id, ad)?; } diff --git a/components/ads-client/src/ffi.rs b/components/ads-client/src/ffi.rs index c17584a1238..5789b0c6c63 100644 --- a/components/ads-client/src/ffi.rs +++ b/components/ads-client/src/ffi.rs @@ -134,6 +134,7 @@ impl MozAdsClientBuilder { .take() .map(MozAdsTelemetryWrapper::new) .unwrap_or_else(MozAdsTelemetryWrapper::noop); + let store_set = inner.store_config.is_some(); let client_config = AdsClientConfig { cache_config: inner.cache_config.clone().map(Into::into), context_id_provider: inner @@ -148,7 +149,11 @@ impl MozAdsClientBuilder { let client = AdsClient::new(client_config); let shutdown_references = client.shutdown_references(); let inner = Arc::new(Mutex::new(client)); - let worker = worker::AdsClientWorkerWrapper::new(inner.clone()); + let worker = if store_set { + worker::AdsClientWorkerWrapper::new(inner.clone()) + } else { + worker::AdsClientWorkerWrapper::new_empty() + }; MozAdsClient { inner, shutdown_references, diff --git a/components/ads-client/src/worker.rs b/components/ads-client/src/worker.rs index 0538d7e99b9..712cec06bf4 100644 --- a/components/ads-client/src/worker.rs +++ b/components/ads-client/src/worker.rs @@ -27,6 +27,13 @@ impl AdsClientWorkerWrapper { } } + pub fn new_empty() -> AdsClientWorkerWrapper { + AdsClientWorkerWrapper { + _worker_thread: None, + worker_dispatch: None, + } + } + pub fn dispatch(&self, command: DispatchCommand) -> Result<(), ComponentError> { if let Some(worker_dispatch) = &self.worker_dispatch { worker_dispatch diff --git a/components/ads-client/src/worker/command.rs b/components/ads-client/src/worker/command.rs index a1bb1ebba4a..de7a12ca53a 100644 --- a/components/ads-client/src/worker/command.rs +++ b/components/ads-client/src/worker/command.rs @@ -1,13 +1,12 @@ use std::{collections::HashMap, sync::mpsc::SyncSender}; use error_support::handle_error; -use url::Url; use crate::{ ads_store::StorableAd, client::error::{BackgroundWorkerError, ComponentError, RequestAdsError}, http_cache::CachePolicy, - mars::{ad_request::AdPlacementRequest, ReportReason}, + mars::ad_request::AdPlacementRequest, AdsClientApiResult, MozAdsClientInner, MozAdsPlacementRequest, MozAdsPlacementRequestWithCount, }; @@ -37,19 +36,6 @@ pub enum DispatchCommand { flags: HashMap, blocks: Vec, }, - RecordClick { - url: Url, - ohttp: bool, - }, - RecordImpression { - url: Url, - ohttp: bool, - }, - ReportAd { - url: Url, - reason: ReportReason, - ohttp: bool, - }, Ping(SyncSender<()>), } @@ -83,7 +69,6 @@ impl DispatchCommand { blocks, ) .map_err(ComponentError::RequestAds)?; - // TODO: Bulk insert inner .cache_ads( image_response @@ -158,27 +143,6 @@ impl DispatchCommand { } Ok(()) } - DispatchCommand::RecordClick { url, ohttp } => { - let inner = ads_client_inner.lock(); - inner - .record_click(url, ohttp) - .map_err(ComponentError::RecordClick)?; - Ok(()) - } - DispatchCommand::RecordImpression { url, ohttp } => { - let inner = ads_client_inner.lock(); - inner - .record_impression(url, ohttp) - .map_err(ComponentError::RecordImpression)?; - Ok(()) - } - DispatchCommand::ReportAd { url, ohttp, reason } => { - let inner = ads_client_inner.lock(); - inner - .report_ad(url, reason, ohttp) - .map_err(ComponentError::ReportAd)?; - Ok(()) - } DispatchCommand::Ping(sender) => { sender .try_send(()) From b187cd4958c400f3c9a1f71d33f55b8ca63587c3 Mon Sep 17 00:00:00 2001 From: Wyatt Verchere Date: Thu, 10 Sep 2026 15:35:41 -0700 Subject: [PATCH 04/12] fix: removed ping --- components/ads-client/src/lib.rs | 17 ----------------- components/ads-client/src/worker/command.rs | 12 +----------- 2 files changed, 1 insertion(+), 28 deletions(-) diff --git a/components/ads-client/src/lib.rs b/components/ads-client/src/lib.rs index f5a0dc71078..3a1e746b2ca 100644 --- a/components/ads-client/src/lib.rs +++ b/components/ads-client/src/lib.rs @@ -191,21 +191,4 @@ impl MozAdsClient { .map_err(ComponentError::RequestAds)?; Ok(response.into_iter().map(|(k, v)| (k, v.into())).collect()) } - - // Pings the background worker and waits for a response back, for use in tests. - // Because the background worker is synchronous, this returns if the worker is empty, - // making it useful for integration tests to wait until all tasks have completed. - #[handle_error(ComponentError)] - pub fn ping_background_worker(&self, timeout: Option) -> AdsClientApiResult<()> { - let (tx, rx) = mpsc::sync_channel(0); - self.worker.dispatch(DispatchCommand::Ping(tx))?; - - if let Some(timeout) = timeout { - rx.recv_timeout(timeout) - .map_err(BackgroundWorkerError::from)?; - } else { - rx.recv().map_err(|_| BackgroundWorkerError::WorkerClosed)?; - } - Ok(()) - } } diff --git a/components/ads-client/src/worker/command.rs b/components/ads-client/src/worker/command.rs index de7a12ca53a..de196407c3b 100644 --- a/components/ads-client/src/worker/command.rs +++ b/components/ads-client/src/worker/command.rs @@ -4,7 +4,7 @@ use error_support::handle_error; use crate::{ ads_store::StorableAd, - client::error::{BackgroundWorkerError, ComponentError, RequestAdsError}, + client::error::{ComponentError, RequestAdsError}, http_cache::CachePolicy, mars::ad_request::AdPlacementRequest, AdsClientApiResult, MozAdsClientInner, MozAdsPlacementRequest, MozAdsPlacementRequestWithCount, @@ -12,8 +12,6 @@ use crate::{ // Command dispatch enum for passing different instructions to the background worker thread. // `RequestImageAds`, `RequestSpocAds`, `RequestTileAds` are prefetch mechanisms that query and load data into the local cache. -// `RecordClick`, `RecordImpression`, `ReportAd` are fire and forget mechanisms that do not load data. -// `Ping` is a synchronous command for internal use that triggers its inner channel when command resolves (eg: when the queue is empty). pub enum DispatchCommand { RequestImageAds { image_ad_requests: Vec, @@ -36,7 +34,6 @@ pub enum DispatchCommand { flags: HashMap, blocks: Vec, }, - Ping(SyncSender<()>), } impl DispatchCommand { @@ -80,7 +77,6 @@ impl DispatchCommand { } Ok(()) } - // TODO: Can we modify these to be one call? DispatchCommand::RequestSpocAds { spoc_ad_requests, cache_policy, @@ -143,12 +139,6 @@ impl DispatchCommand { } Ok(()) } - DispatchCommand::Ping(sender) => { - sender - .try_send(()) - .map_err(|err| BackgroundWorkerError::PongFailure(Box::new(err)))?; - Ok(()) - } } } } From 0e4c7d53dbd43a863fc680b8b19e1e9e7a91796e Mon Sep 17 00:00:00 2001 From: Wyatt Verchere Date: Thu, 10 Sep 2026 16:00:33 -0700 Subject: [PATCH 05/12] fix: pong removal --- components/ads-client/src/client/error.rs | 3 --- components/ads-client/src/ffi.rs | 2 +- components/ads-client/src/lib.rs | 8 +++----- components/ads-client/src/worker/command.rs | 6 +----- 4 files changed, 5 insertions(+), 14 deletions(-) diff --git a/components/ads-client/src/client/error.rs b/components/ads-client/src/client/error.rs index 913a36792c9..cf5371653a6 100644 --- a/components/ads-client/src/client/error.rs +++ b/components/ads-client/src/client/error.rs @@ -39,9 +39,6 @@ pub enum RequestAdsError { #[derive(Debug, thiserror::Error)] pub enum BackgroundWorkerError { - #[error("Error sending pong back from background worker")] - PongFailure(Box>), - #[error("Error requesting new ads from the background worker: worker closed")] WorkerClosed, diff --git a/components/ads-client/src/ffi.rs b/components/ads-client/src/ffi.rs index 5789b0c6c63..49b462515a4 100644 --- a/components/ads-client/src/ffi.rs +++ b/components/ads-client/src/ffi.rs @@ -157,7 +157,7 @@ impl MozAdsClientBuilder { MozAdsClient { inner, shutdown_references, - worker, + _worker: worker } } diff --git a/components/ads-client/src/lib.rs b/components/ads-client/src/lib.rs index 3a1e746b2ca..7d387d82ff9 100644 --- a/components/ads-client/src/lib.rs +++ b/components/ads-client/src/lib.rs @@ -5,8 +5,7 @@ use std::{ collections::HashMap, - sync::{mpsc, Arc}, - time::Duration, + sync::Arc, }; use client::error::ComponentError; @@ -31,10 +30,9 @@ pub mod worker; pub use ffi::*; use crate::{ - client::error::BackgroundWorkerError, ffi::telemetry::MozAdsTelemetryWrapper, shutdown::ShutdownReferences, - worker::{command::DispatchCommand, AdsClientWorkerWrapper}, + worker::AdsClientWorkerWrapper, }; #[cfg(test)] @@ -53,7 +51,7 @@ pub type MozAdsClientInner = Arc>>; pub struct MozAdsClient { inner: MozAdsClientInner, shutdown_references: ShutdownReferences, - worker: AdsClientWorkerWrapper, + _worker: AdsClientWorkerWrapper, } #[uniffi::export] diff --git a/components/ads-client/src/worker/command.rs b/components/ads-client/src/worker/command.rs index de196407c3b..a443fe2ce6b 100644 --- a/components/ads-client/src/worker/command.rs +++ b/components/ads-client/src/worker/command.rs @@ -1,4 +1,4 @@ -use std::{collections::HashMap, sync::mpsc::SyncSender}; +use std::collections::HashMap; use error_support::handle_error; @@ -39,10 +39,6 @@ pub enum DispatchCommand { impl DispatchCommand { // Runs a dispatched command synchronously in it's thread. // The dispatched command calls the corresponding `AdsClient` synchronous method, meaning that behavior between the two is shared. - // This includes telemetry calls, meaning that for a successful `RecordClick`, all of the following will get logged: - // - CommandDispatchedOperationEvent::RecordClick (on dispatch) - // - ClientOperationEvent::RecordClick (on `AdsClient` method success) - // - CommandProcessedOperationEvent::RecordClick (on process) #[handle_error(ComponentError)] pub fn run_command(self, ads_client_inner: &MozAdsClientInner) -> AdsClientApiResult<()> { match self { From 59fa8318a23e86f1872e7e9ac3d2c3079e10ac81 Mon Sep 17 00:00:00 2001 From: Wyatt Verchere Date: Thu, 10 Sep 2026 16:36:04 -0700 Subject: [PATCH 06/12] fix: adds shutdown --- components/ads-client/src/ads_store.rs | 27 ----------------------- components/ads-client/src/client.rs | 8 +++---- components/ads-client/src/client/error.rs | 12 +++++----- components/ads-client/src/ffi.rs | 2 +- components/ads-client/src/lib.rs | 8 ++----- components/ads-client/src/mars/error.rs | 3 +++ components/ads-client/src/worker.rs | 2 +- 7 files changed, 17 insertions(+), 45 deletions(-) diff --git a/components/ads-client/src/ads_store.rs b/components/ads-client/src/ads_store.rs index bf72b36982d..d8db0b7a5d1 100644 --- a/components/ads-client/src/ads_store.rs +++ b/components/ads-client/src/ads_store.rs @@ -49,33 +49,6 @@ pub enum StorableAd { Tile(AdTile), } -impl StorableAd { - // TODO: Should these be result? - pub fn into_image(self) -> Option { - if let StorableAd::Image(image) = self { - Some(image) - } else { - None - } - } - - pub fn into_spocs(self) -> Option> { - if let StorableAd::Spoc(spocs) = self { - Some(spocs) - } else { - None - } - } - - pub fn into_tile(self) -> Option { - if let StorableAd::Tile(tile) = self { - Some(tile) - } else { - None - } - } -} - pub struct AdsStore { holder: AdsStoreHolder, #[allow(dead_code)] diff --git a/components/ads-client/src/client.rs b/components/ads-client/src/client.rs index 739892f6fc5..afad1f43e13 100644 --- a/components/ads-client/src/client.rs +++ b/components/ads-client/src/client.rs @@ -123,9 +123,10 @@ where let ads_store = self.ads_store.lock(); if let Some(ads_store) = ads_store.as_ref() { ads_store.store_ads(ads)?; + Ok(()) + } else { + Err(FetchAdsError::SqliteShutdown) } - // TODO: Should we error if no ads store? - Ok(()) } pub fn get_cached_ad( @@ -136,8 +137,7 @@ where if let Some(ads_store) = ads_store.as_ref() { Ok(ads_store.lookup(placement_id)?) } else { - // TODO: Should we error if no ads store? - Ok(None) + Err(FetchAdsError::SqliteShutdown) } } diff --git a/components/ads-client/src/client/error.rs b/components/ads-client/src/client/error.rs index cf5371653a6..ea1b96f24a7 100644 --- a/components/ads-client/src/client/error.rs +++ b/components/ads-client/src/client/error.rs @@ -40,20 +40,20 @@ pub enum RequestAdsError { #[derive(Debug, thiserror::Error)] pub enum BackgroundWorkerError { #[error("Error requesting new ads from the background worker: worker closed")] - WorkerClosed, + Closed, #[error("Error requesting new ads from the background worker: worker full")] - WorkerFull, + Full, - #[error("Worker timed out waiting for response: {0}")] - WorkerTimedOut(#[from] RecvTimeoutError), + #[error("Background worker timed out waiting for response: {0}")] + TimedOut(#[from] RecvTimeoutError), } impl From> for BackgroundWorkerError { fn from(value: TrySendError) -> Self { match value { - TrySendError::Disconnected(_) => BackgroundWorkerError::WorkerClosed, - TrySendError::Full(_) => BackgroundWorkerError::WorkerFull, + TrySendError::Disconnected(_) => BackgroundWorkerError::Closed, + TrySendError::Full(_) => BackgroundWorkerError::Full, } } } diff --git a/components/ads-client/src/ffi.rs b/components/ads-client/src/ffi.rs index 49b462515a4..bfec99326e9 100644 --- a/components/ads-client/src/ffi.rs +++ b/components/ads-client/src/ffi.rs @@ -157,7 +157,7 @@ impl MozAdsClientBuilder { MozAdsClient { inner, shutdown_references, - _worker: worker + _worker: worker, } } diff --git a/components/ads-client/src/lib.rs b/components/ads-client/src/lib.rs index 7d387d82ff9..11dc7789d2b 100644 --- a/components/ads-client/src/lib.rs +++ b/components/ads-client/src/lib.rs @@ -3,10 +3,7 @@ * file, You can obtain one at http://mozilla.org/MPL/2.0/. */ -use std::{ - collections::HashMap, - sync::Arc, -}; +use std::{collections::HashMap, sync::Arc}; use client::error::ComponentError; use error_support::handle_error; @@ -30,8 +27,7 @@ pub mod worker; pub use ffi::*; use crate::{ - ffi::telemetry::MozAdsTelemetryWrapper, - shutdown::ShutdownReferences, + ffi::telemetry::MozAdsTelemetryWrapper, shutdown::ShutdownReferences, worker::AdsClientWorkerWrapper, }; diff --git a/components/ads-client/src/mars/error.rs b/components/ads-client/src/mars/error.rs index ce87833924c..22f9f0573dd 100644 --- a/components/ads-client/src/mars/error.rs +++ b/components/ads-client/src/mars/error.rs @@ -46,6 +46,9 @@ pub enum FetchAdsError { #[error("Internal database error: {0}")] Sqlite(#[from] rusqlite::Error), + #[error("Internal database error: database shut down or uninitialized")] + SqliteShutdown, + #[error("Error sending request: {0}")] Request(#[from] viaduct::ViaductError), diff --git a/components/ads-client/src/worker.rs b/components/ads-client/src/worker.rs index 712cec06bf4..6af30746715 100644 --- a/components/ads-client/src/worker.rs +++ b/components/ads-client/src/worker.rs @@ -42,7 +42,7 @@ impl AdsClientWorkerWrapper { Ok(()) } else { - Err(BackgroundWorkerError::WorkerClosed.into()) + Err(BackgroundWorkerError::Closed.into()) } } } From 93f7d7dbf8e696195f50b2532317ad424fc1f0b7 Mon Sep 17 00:00:00 2001 From: Wyatt Verchere Date: Thu, 10 Sep 2026 17:29:17 -0700 Subject: [PATCH 07/12] fix: buffer --- components/ads-client/src/ffi.rs | 7 ++++++- components/ads-client/src/worker.rs | 17 +++++++++++++---- 2 files changed, 19 insertions(+), 5 deletions(-) diff --git a/components/ads-client/src/ffi.rs b/components/ads-client/src/ffi.rs index bfec99326e9..924991ab3a3 100644 --- a/components/ads-client/src/ffi.rs +++ b/components/ads-client/src/ffi.rs @@ -135,6 +135,10 @@ impl MozAdsClientBuilder { .map(MozAdsTelemetryWrapper::new) .unwrap_or_else(MozAdsTelemetryWrapper::noop); let store_set = inner.store_config.is_some(); + let worker_buffer_size = inner + .store_config + .as_ref() + .and_then(|x| x.worker_buffer_size); let client_config = AdsClientConfig { cache_config: inner.cache_config.clone().map(Into::into), context_id_provider: inner @@ -150,7 +154,7 @@ impl MozAdsClientBuilder { let shutdown_references = client.shutdown_references(); let inner = Arc::new(Mutex::new(client)); let worker = if store_set { - worker::AdsClientWorkerWrapper::new(inner.clone()) + worker::AdsClientWorkerWrapper::new(inner.clone(), worker_buffer_size) } else { worker::AdsClientWorkerWrapper::new_empty() }; @@ -219,6 +223,7 @@ pub struct MozAdsCacheConfig { #[derive(Clone, uniffi::Record)] pub struct MozAdsStoreConfig { pub db_path: String, + pub worker_buffer_size: Option, } #[derive(Debug, PartialEq, uniffi::Record)] diff --git a/components/ads-client/src/worker.rs b/components/ads-client/src/worker.rs index 6af30746715..d01ef1e2065 100644 --- a/components/ads-client/src/worker.rs +++ b/components/ads-client/src/worker.rs @@ -10,7 +10,8 @@ use std::{ pub mod command; -pub const ADS_CLIENT_WORKER_CHANNEL_BUFFER_SIZE: usize = 1000; +// This is a somewhat arbitrary default value that is overridable. +pub const ADS_CLIENT_WORKER_CHANNEL_BUFFER_SIZE_DEFAULT: usize = 10000; pub const ADS_CLIENT_WORKER_THREAD_NAME: &str = "ads-client.worker"; pub struct AdsClientWorkerWrapper { @@ -19,8 +20,13 @@ pub struct AdsClientWorkerWrapper { } impl AdsClientWorkerWrapper { - pub fn new(inner: MozAdsClientInner) -> AdsClientWorkerWrapper { - let (worker_dispatch, worker_thread) = Option::unzip(build_worker_thread(inner.clone())); + pub fn new( + inner: MozAdsClientInner, + worker_buffer_size: Option, + ) -> AdsClientWorkerWrapper { + let worker_buffer_size = worker_buffer_size.and_then(|x| usize::try_from(x).ok()); + let (worker_dispatch, worker_thread) = + Option::unzip(build_worker_thread(inner.clone(), worker_buffer_size)); AdsClientWorkerWrapper { _worker_thread: worker_thread, worker_dispatch, @@ -51,8 +57,11 @@ impl AdsClientWorkerWrapper { // Returns None if thread fails to build. pub fn build_worker_thread( inner_client: MozAdsClientInner, + max_channel_size: Option, ) -> Option<(SyncSender, JoinHandle<()>)> { - let (tx, rx) = mpsc::sync_channel(ADS_CLIENT_WORKER_CHANNEL_BUFFER_SIZE); + let (tx, rx) = mpsc::sync_channel( + max_channel_size.unwrap_or(ADS_CLIENT_WORKER_CHANNEL_BUFFER_SIZE_DEFAULT), + ); let worker_thread_handle = std::thread::Builder::new() .name(ADS_CLIENT_WORKER_THREAD_NAME.to_string()) .spawn(move || crate::worker::worker(inner_client, rx)).inspect_err(|err| { From 231d4cddae3ae900eb1bee4d631b11166fc901b0 Mon Sep 17 00:00:00 2001 From: Wyatt Verchere Date: Thu, 10 Sep 2026 17:40:20 -0700 Subject: [PATCH 08/12] fix: changelog --- CHANGELOG.md | 6 ++++++ components/ads-client/src/worker.rs | 2 +- 2 files changed, 7 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 4850659024f..6e6aba952b7 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,12 @@ [Full Changelog](In progress) +## ✨ What's Changed ✨ + +### Ads-Client + +- Adds a background worker to allow for future fire-and-forget logic. + # v157.0 (_2026-09-10_) ## ✨ What's Changed ✨ diff --git a/components/ads-client/src/worker.rs b/components/ads-client/src/worker.rs index d01ef1e2065..039c84e2290 100644 --- a/components/ads-client/src/worker.rs +++ b/components/ads-client/src/worker.rs @@ -10,7 +10,7 @@ use std::{ pub mod command; -// This is a somewhat arbitrary default value that is overridable. +// This is a somewhat arbitrary default value that is overridable. pub const ADS_CLIENT_WORKER_CHANNEL_BUFFER_SIZE_DEFAULT: usize = 10000; pub const ADS_CLIENT_WORKER_THREAD_NAME: &str = "ads-client.worker"; From df0fcac0f578a72a83f7b0dfd5ac849d63bf5100 Mon Sep 17 00:00:00 2001 From: Wyatt Verchere Date: Fri, 11 Sep 2026 17:45:26 -0700 Subject: [PATCH 09/12] fix: Revisions --- CHANGELOG.md | 6 --- components/ads-client/src/ads_store.rs | 26 +++++++++++ components/ads-client/src/client.rs | 50 ++++++++++++++++++--- components/ads-client/src/ffi.rs | 4 +- components/ads-client/src/lib.rs | 5 +-- components/ads-client/src/worker.rs | 15 +++---- components/ads-client/src/worker/command.rs | 6 +-- 7 files changed, 82 insertions(+), 30 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index ca50530a873..77ae367bc02 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,12 +2,6 @@ [Full Changelog](In progress) -## ✨ What's Changed ✨ - -### Ads-Client - -- Adds a background worker to allow for future fire-and-forget logic. - ### Glean - Updated to v70.0.0 ([#7598](https://github.com/mozilla/application-services/pull/7598)) diff --git a/components/ads-client/src/ads_store.rs b/components/ads-client/src/ads_store.rs index d8db0b7a5d1..17ecf9fa891 100644 --- a/components/ads-client/src/ads_store.rs +++ b/components/ads-client/src/ads_store.rs @@ -49,6 +49,32 @@ pub enum StorableAd { Tile(AdTile), } +impl StorableAd { + pub fn into_image(self) -> Option { + if let StorableAd::Image(image) = self { + Some(image) + } else { + None + } + } + + pub fn into_spocs(self) -> Option> { + if let StorableAd::Spoc(spocs) = self { + Some(spocs) + } else { + None + } + } + + pub fn into_tile(self) -> Option { + if let StorableAd::Tile(tile) = self { + Some(tile) + } else { + None + } + } +} + pub struct AdsStore { holder: AdsStoreHolder, #[allow(dead_code)] diff --git a/components/ads-client/src/client.rs b/components/ads-client/src/client.rs index afad1f43e13..94a9d91d31f 100644 --- a/components/ads-client/src/client.rs +++ b/components/ads-client/src/client.rs @@ -116,7 +116,7 @@ where self.client.clear_cache() } - pub fn cache_ads( + pub fn store_ads( &mut self, ads: HashMap, ) -> Result<(), FetchAdsError> { @@ -129,15 +129,51 @@ where } } - pub fn get_cached_ad( - &self, - placement_id: &PlacementId, - ) -> Result, FetchAdsError> { + pub fn get_stored_ad_images(&self, placement_id: &PlacementId) -> Option { let ads_store = self.ads_store.lock(); if let Some(ads_store) = ads_store.as_ref() { - Ok(ads_store.lookup(placement_id)?) + match ads_store.lookup(placement_id) { + Ok(ad) => ad.and_then(|ad| ad.into_image()), + Err(_) => { + // TODO: Telemetry should return an error here (eg: some internal sqlite error) + None + } + } } else { - Err(FetchAdsError::SqliteShutdown) + // TODO: Telemetry should be added here for the database being shut down. + None + } + } + + pub fn get_stored_ad_spocs(&self, placement_id: &PlacementId) -> Option> { + let ads_store = self.ads_store.lock(); + if let Some(ads_store) = ads_store.as_ref() { + match ads_store.lookup(placement_id) { + Ok(ad) => ad.and_then(|ad| ad.into_spocs()), + Err(_) => { + // TODO: Telemetry should return an error here (eg: some internal sqlite error) + None + } + } + } else { + // TODO: Telemetry should be added here for the database being shut down. + None + } + } + + pub fn get_stored_ad_tile(&self, placement_id: &PlacementId) -> Option { + let ads_store = self.ads_store.lock(); + if let Some(ads_store) = ads_store.as_ref() { + match ads_store.lookup(placement_id) { + Ok(ad) => ad.and_then(|ad| ad.into_tile()), + Err(_) => { + // TODO: Telemetry should return an error here (eg: some internal sqlite error) + None + } + } + } else { + // TODO: Telemetry should be added here for the database being shut down. + None } } diff --git a/components/ads-client/src/ffi.rs b/components/ads-client/src/ffi.rs index 924991ab3a3..8e8120069b8 100644 --- a/components/ads-client/src/ffi.rs +++ b/components/ads-client/src/ffi.rs @@ -154,9 +154,9 @@ impl MozAdsClientBuilder { let shutdown_references = client.shutdown_references(); let inner = Arc::new(Mutex::new(client)); let worker = if store_set { - worker::AdsClientWorkerWrapper::new(inner.clone(), worker_buffer_size) + worker::BackgroundWorker::new(inner.clone(), worker_buffer_size) } else { - worker::AdsClientWorkerWrapper::new_empty() + worker::BackgroundWorker::new_empty() }; MozAdsClient { inner, diff --git a/components/ads-client/src/lib.rs b/components/ads-client/src/lib.rs index 11dc7789d2b..a1141d63d63 100644 --- a/components/ads-client/src/lib.rs +++ b/components/ads-client/src/lib.rs @@ -27,8 +27,7 @@ pub mod worker; pub use ffi::*; use crate::{ - ffi::telemetry::MozAdsTelemetryWrapper, shutdown::ShutdownReferences, - worker::AdsClientWorkerWrapper, + ffi::telemetry::MozAdsTelemetryWrapper, shutdown::ShutdownReferences, worker::BackgroundWorker, }; #[cfg(test)] @@ -47,7 +46,7 @@ pub type MozAdsClientInner = Arc>>; pub struct MozAdsClient { inner: MozAdsClientInner, shutdown_references: ShutdownReferences, - _worker: AdsClientWorkerWrapper, + _worker: BackgroundWorker, } #[uniffi::export] diff --git a/components/ads-client/src/worker.rs b/components/ads-client/src/worker.rs index 039c84e2290..1de19b25d4f 100644 --- a/components/ads-client/src/worker.rs +++ b/components/ads-client/src/worker.rs @@ -14,27 +14,24 @@ pub mod command; pub const ADS_CLIENT_WORKER_CHANNEL_BUFFER_SIZE_DEFAULT: usize = 10000; pub const ADS_CLIENT_WORKER_THREAD_NAME: &str = "ads-client.worker"; -pub struct AdsClientWorkerWrapper { +pub struct BackgroundWorker { _worker_thread: Option>, worker_dispatch: Option>, } -impl AdsClientWorkerWrapper { - pub fn new( - inner: MozAdsClientInner, - worker_buffer_size: Option, - ) -> AdsClientWorkerWrapper { +impl BackgroundWorker { + pub fn new(inner: MozAdsClientInner, worker_buffer_size: Option) -> BackgroundWorker { let worker_buffer_size = worker_buffer_size.and_then(|x| usize::try_from(x).ok()); let (worker_dispatch, worker_thread) = Option::unzip(build_worker_thread(inner.clone(), worker_buffer_size)); - AdsClientWorkerWrapper { + BackgroundWorker { _worker_thread: worker_thread, worker_dispatch, } } - pub fn new_empty() -> AdsClientWorkerWrapper { - AdsClientWorkerWrapper { + pub fn new_empty() -> BackgroundWorker { + BackgroundWorker { _worker_thread: None, worker_dispatch: None, } diff --git a/components/ads-client/src/worker/command.rs b/components/ads-client/src/worker/command.rs index a443fe2ce6b..33783a2d6dd 100644 --- a/components/ads-client/src/worker/command.rs +++ b/components/ads-client/src/worker/command.rs @@ -63,7 +63,7 @@ impl DispatchCommand { ) .map_err(ComponentError::RequestAds)?; inner - .cache_ads( + .store_ads( image_response .into_iter() .map(|(k, v)| (k.into(), StorableAd::Image(v))) @@ -94,7 +94,7 @@ impl DispatchCommand { ) .map_err(ComponentError::RequestAds)?; inner - .cache_ads( + .store_ads( spoc_response .into_iter() .map(|(k, v)| (k.into(), StorableAd::Spoc(v))) @@ -125,7 +125,7 @@ impl DispatchCommand { ) .map_err(ComponentError::RequestAds)?; inner - .cache_ads( + .store_ads( tile_response .into_iter() .map(|(k, v)| (k.into(), StorableAd::Tile(v))) From 141f1d50c1b7cce11dd48f5beb0573bbe95b9509 Mon Sep 17 00:00:00 2001 From: Wyatt Verchere Date: Tue, 15 Sep 2026 15:10:29 -0700 Subject: [PATCH 10/12] fix: updated with cfg --- components/ads-client/src/client.rs | 8 +++++++- components/ads-client/src/client/error.rs | 12 ++++++------ components/ads-client/src/ffi.rs | 8 +++++++- components/ads-client/src/lib.rs | 8 +++++--- 4 files changed, 25 insertions(+), 11 deletions(-) diff --git a/components/ads-client/src/client.rs b/components/ads-client/src/client.rs index da09ade8fe3..8c453fedbe3 100644 --- a/components/ads-client/src/client.rs +++ b/components/ads-client/src/client.rs @@ -9,7 +9,9 @@ use crate::common::bytesize::ByteSize; use crate::http_cache::{CachePolicy, HttpCache}; use crate::mars::ad_request::{AdPlacementRequest, AdRequestFlags}; use crate::mars::ad_response::{AdImage, AdResponse, AdResponseValue, AdSpoc, AdTile}; -use crate::mars::error::{FetchAdsError, RecordClickError, RecordImpressionError, ReportAdError}; +#[cfg(feature = "stateful")] +use crate::mars::error::FetchAdsError; +use crate::mars::error::{RecordClickError, RecordImpressionError, ReportAdError}; use crate::mars::{MARSClient, ReportReason}; #[cfg(feature = "stateful")] use crate::shutdown::AdsStoreShutdown; @@ -111,6 +113,7 @@ where self.client.clear_cache() } + #[cfg(feature = "stateful")] pub fn store_ads( &mut self, ads: HashMap, @@ -124,6 +127,7 @@ where } } + #[cfg(feature = "stateful")] pub fn get_stored_ad_images(&self, placement_id: &PlacementId) -> Option { let ads_store = self.ads_store.lock(); if let Some(ads_store) = ads_store.as_ref() { @@ -140,6 +144,7 @@ where } } + #[cfg(feature = "stateful")] pub fn get_stored_ad_spocs(&self, placement_id: &PlacementId) -> Option> { let ads_store = self.ads_store.lock(); if let Some(ads_store) = ads_store.as_ref() { @@ -156,6 +161,7 @@ where } } + #[cfg(feature = "stateful")] pub fn get_stored_ad_tile(&self, placement_id: &PlacementId) -> Option { let ads_store = self.ads_store.lock(); if let Some(ads_store) = ads_store.as_ref() { diff --git a/components/ads-client/src/client/error.rs b/components/ads-client/src/client/error.rs index ea1b96f24a7..9067f1bf9c5 100644 --- a/components/ads-client/src/client/error.rs +++ b/components/ads-client/src/client/error.rs @@ -3,15 +3,13 @@ * file, You can obtain one at http://mozilla.org/MPL/2.0/. */ -use std::sync::mpsc::{RecvTimeoutError, TrySendError}; - -use crate::{ - mars::error::{FetchAdsError, RecordClickError, RecordImpressionError, ReportAdError}, - worker::command, -}; +use crate::mars::error::{FetchAdsError, RecordClickError, RecordImpressionError, ReportAdError}; +#[cfg(feature = "stateful")] +use crate::worker::command; #[derive(Debug, thiserror::Error)] pub enum ComponentError { + #[cfg(feature = "stateful")] #[error("Error requesting ads from worker: {0}")] BackgroundWorker(#[from] BackgroundWorkerError), @@ -37,6 +35,7 @@ pub enum RequestAdsError { FetchAds(#[from] FetchAdsError), } +#[cfg(feature = "stateful")] #[derive(Debug, thiserror::Error)] pub enum BackgroundWorkerError { #[error("Error requesting new ads from the background worker: worker closed")] @@ -49,6 +48,7 @@ pub enum BackgroundWorkerError { TimedOut(#[from] RecvTimeoutError), } +#[cfg(feature = "stateful")] impl From> for BackgroundWorkerError { fn from(value: TrySendError) -> Self { match value { diff --git a/components/ads-client/src/ffi.rs b/components/ads-client/src/ffi.rs index d956be28bbb..e75e1ff5744 100644 --- a/components/ads-client/src/ffi.rs +++ b/components/ads-client/src/ffi.rs @@ -20,8 +20,10 @@ use crate::mars::ad_response::{ }; use crate::mars::Environment; use crate::mars::ReportReason; +#[cfg(feature = "stateful")] +use crate::worker; use crate::AdsClientUrl; -use crate::{worker, MozAdsClient}; +use crate::MozAdsClient; use parking_lot::Mutex; use std::collections::HashMap; use std::sync::Arc; @@ -107,7 +109,9 @@ impl MozAdsClientBuilder { .take() .map(MozAdsTelemetryWrapper::new) .unwrap_or_else(MozAdsTelemetryWrapper::noop); + #[cfg(feature = "stateful")] let store_set = inner.store_config.is_some(); + #[cfg(feature = "stateful")] let worker_buffer_size = inner .store_config .as_ref() @@ -122,6 +126,7 @@ impl MozAdsClientBuilder { let client = AdsClient::new(client_config); let shutdown_references = client.shutdown_references(); let inner = Arc::new(Mutex::new(client)); + #[cfg(feature = "stateful")] let worker = if store_set { worker::BackgroundWorker::new(inner.clone(), worker_buffer_size) } else { @@ -130,6 +135,7 @@ impl MozAdsClientBuilder { MozAdsClient { inner, shutdown_references, + #[cfg(feature = "stateful")] _worker: worker, } } diff --git a/components/ads-client/src/lib.rs b/components/ads-client/src/lib.rs index 733a7705385..587fd3616a5 100644 --- a/components/ads-client/src/lib.rs +++ b/components/ads-client/src/lib.rs @@ -23,13 +23,14 @@ pub mod http_cache; mod mars; pub mod shutdown; pub mod telemetry; +#[cfg(feature = "stateful")] pub mod worker; pub use ffi::*; -use crate::{ - ffi::telemetry::MozAdsTelemetryWrapper, shutdown::ShutdownReferences, worker::BackgroundWorker, -}; +#[cfg(feature = "stateful")] +use crate::worker::BackgroundWorker; +use crate::{ffi::telemetry::MozAdsTelemetryWrapper, shutdown::ShutdownReferences}; #[cfg(test)] mod test_utils; @@ -47,6 +48,7 @@ pub type MozAdsClientInner = Arc>>; pub struct MozAdsClient { inner: MozAdsClientInner, shutdown_references: ShutdownReferences, + #[cfg(feature = "stateful")] _worker: BackgroundWorker, } From 46b0abe29c414a07755cf2d2648f21812f9c63b8 Mon Sep 17 00:00:00 2001 From: Wyatt Verchere Date: Tue, 15 Sep 2026 15:13:49 -0700 Subject: [PATCH 11/12] fix: fmt clippy --- components/ads-client/src/client/error.rs | 3 +++ 1 file changed, 3 insertions(+) diff --git a/components/ads-client/src/client/error.rs b/components/ads-client/src/client/error.rs index 9067f1bf9c5..e6e648ce9c1 100644 --- a/components/ads-client/src/client/error.rs +++ b/components/ads-client/src/client/error.rs @@ -3,6 +3,9 @@ * file, You can obtain one at http://mozilla.org/MPL/2.0/. */ +#[cfg(feature = "stateful")] +use std::sync::mpsc::{RecvTimeoutError, TrySendError}; + use crate::mars::error::{FetchAdsError, RecordClickError, RecordImpressionError, ReportAdError}; #[cfg(feature = "stateful")] use crate::worker::command; From d02d8a136cf2fa5298172ad39b92cd6e4a9ae6b3 Mon Sep 17 00:00:00 2001 From: Wyatt Verchere Date: Thu, 17 Sep 2026 10:09:17 -0700 Subject: [PATCH 12/12] fix: reorganizes background worker part --- components/ads-client/src/worker.rs | 24 +++++++++++++++++++----- 1 file changed, 19 insertions(+), 5 deletions(-) diff --git a/components/ads-client/src/worker.rs b/components/ads-client/src/worker.rs index 1de19b25d4f..e679d572cae 100644 --- a/components/ads-client/src/worker.rs +++ b/components/ads-client/src/worker.rs @@ -20,13 +20,27 @@ pub struct BackgroundWorker { } impl BackgroundWorker { - pub fn new(inner: MozAdsClientInner, worker_buffer_size: Option) -> BackgroundWorker { + pub fn new( + inner_client: MozAdsClientInner, + worker_buffer_size: Option, + ) -> BackgroundWorker { let worker_buffer_size = worker_buffer_size.and_then(|x| usize::try_from(x).ok()); - let (worker_dispatch, worker_thread) = - Option::unzip(build_worker_thread(inner.clone(), worker_buffer_size)); + + let (tx, rx) = mpsc::sync_channel( + worker_buffer_size.unwrap_or(ADS_CLIENT_WORKER_CHANNEL_BUFFER_SIZE_DEFAULT), + ); + let Some(worker_thread) = std::thread::Builder::new() + .name(ADS_CLIENT_WORKER_THREAD_NAME.to_string()) + .spawn(move || crate::worker::worker(inner_client, rx)).inspect_err(|err| { + error_support::error!("Failed to create ads-client worker thread `{ADS_CLIENT_WORKER_THREAD_NAME}` with: {err}") + }).ok() + else { + return BackgroundWorker { _worker_thread: None, worker_dispatch: None } + }; + BackgroundWorker { - _worker_thread: worker_thread, - worker_dispatch, + _worker_thread: Some(worker_thread), + worker_dispatch: Some(tx), } }