use crate::metrics::{registry, Result};
use crate::sdk::{
export::metrics::{AggregatorSelector, CheckpointSet, Checkpointer, ExportKindFor, Record},
metrics::{
accumulator,
processors::{self, BasicProcessor},
Accumulator,
},
Resource,
};
use std::sync::Arc;
use std::time::{Duration, SystemTime};
const DEFAULT_CACHE_DURATION: Duration = Duration::from_secs(10);
pub fn pull(
aggregator_selector: Box<dyn AggregatorSelector + Send + Sync>,
export_selector: Box<dyn ExportKindFor + Send + Sync>,
) -> PullControllerBuilder {
PullControllerBuilder::with_selectors(aggregator_selector, export_selector)
}
#[derive(Debug)]
pub struct PullController {
accumulator: Accumulator,
processor: Arc<BasicProcessor>,
provider: registry::RegistryMeterProvider,
period: Duration,
last_collect: SystemTime,
}
impl PullController {
pub fn provider(&self) -> registry::RegistryMeterProvider {
self.provider.clone()
}
pub fn collect(&mut self) -> Result<()> {
if self
.last_collect
.elapsed()
.map_or(true, |elapsed| elapsed > self.period)
{
self.last_collect = crate::time::now();
self.processor.lock().and_then(|mut checkpointer| {
checkpointer.start_collection();
self.accumulator.0.collect(&mut checkpointer);
checkpointer.finish_collection()
})
} else {
Ok(())
}
}
}
impl CheckpointSet for PullController {
fn try_for_each(
&mut self,
export_selector: &dyn ExportKindFor,
f: &mut dyn FnMut(&Record<'_>) -> Result<()>,
) -> Result<()> {
self.processor.lock().and_then(|mut locked_processor| {
locked_processor
.checkpoint_set()
.try_for_each(export_selector, f)
})
}
}
#[derive(Debug)]
pub struct PullControllerBuilder {
aggregator_selector: Box<dyn AggregatorSelector + Send + Sync>,
export_selector: Box<dyn ExportKindFor + Send + Sync>,
resource: Option<Resource>,
cache_period: Option<Duration>,
memory: bool,
}
impl PullControllerBuilder {
pub fn with_selectors(
aggregator_selector: Box<dyn AggregatorSelector + Send + Sync>,
export_selector: Box<dyn ExportKindFor + Send + Sync>,
) -> Self {
PullControllerBuilder {
aggregator_selector,
export_selector,
resource: None,
cache_period: None,
memory: true,
}
}
pub fn with_resource(self, resource: Resource) -> Self {
PullControllerBuilder {
resource: Some(resource),
..self
}
}
pub fn with_cache_period(self, period: Duration) -> Self {
PullControllerBuilder {
cache_period: Some(period),
..self
}
}
pub fn with_memory(self, memory: bool) -> Self {
PullControllerBuilder { memory, ..self }
}
pub fn build(self) -> PullController {
let processor = Arc::new(processors::basic(
self.aggregator_selector,
self.export_selector,
self.memory,
));
let accumulator = accumulator(processor.clone())
.with_resource(self.resource.unwrap_or_default())
.build();
let provider = registry::meter_provider(Arc::new(accumulator.clone()));
PullController {
accumulator,
processor,
provider,
period: self.cache_period.unwrap_or(DEFAULT_CACHE_DURATION),
last_collect: crate::time::now(),
}
}
}