use std::collections::BTreeMap;
use std::io::Write;
use crate::error::{JobError, Result};
const STATISTICS_FD: i32 = 5;
const MAX_STATISTICS: usize = 128;
#[derive(Debug, Default)]
pub struct JobStatistics {
values: BTreeMap<String, i64>,
finished: bool,
}
impl JobStatistics {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn set(&mut self, name: impl Into<String>, value: i64) -> Result<()> {
let name = name.into();
self.reserve(&name)?;
self.values.insert(name, value);
Ok(())
}
pub fn add(&mut self, name: impl Into<String>, delta: i64) -> Result<()> {
let name = name.into();
self.reserve(&name)?;
let slot = self.values.entry(name).or_insert(0);
*slot = slot.saturating_add(delta);
Ok(())
}
#[must_use]
pub fn get(&self, name: &str) -> Option<i64> {
self.values.get(name).copied()
}
#[must_use]
pub fn len(&self) -> usize {
self.values.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.values.is_empty()
}
pub fn finish(&mut self) -> Result<()> {
if self.finished || self.values.is_empty() {
return Ok(());
}
self.finished = true;
if !crate::is_inside_job() {
eprintln!(
"ytsaurus-job: not running as a job, so {} statistic(s) were not sent",
self.values.len()
);
return Ok(());
}
let encoded = self.encode()?;
write_to_statistics_fd(&encoded)
}
fn encode(&self) -> Result<Vec<u8>> {
let map = ytsaurus_yson::YsonValue {
attributes: None,
node: ytsaurus_yson::YsonNode::Map(
self.values
.iter()
.map(|(name, value)| {
(
name.as_bytes().to_vec(),
ytsaurus_yson::YsonValue {
attributes: None,
node: ytsaurus_yson::YsonNode::Int64(*value),
},
)
})
.collect(),
),
};
let mut encoded =
ytsaurus_yson::to_vec(&map, ytsaurus_yson::YsonFormat::Text).map_err(|e| {
JobError::Statistics {
reason: format!("could not encode them: {e}"),
}
})?;
encoded.push(b';');
Ok(encoded)
}
fn reserve(&self, name: &str) -> Result<()> {
if self.values.len() >= MAX_STATISTICS && !self.values.contains_key(name) {
return Err(JobError::TooManyStatistics {
limit: MAX_STATISTICS,
name: name.to_owned(),
});
}
Ok(())
}
}
impl Drop for JobStatistics {
fn drop(&mut self) {
if self.finished || self.values.is_empty() {
return;
}
if let Err(e) = self.finish() {
eprintln!("ytsaurus-job: could not send job statistics: {e}");
}
}
}
fn write_to_statistics_fd(bytes: &[u8]) -> Result<()> {
use std::mem::ManuallyDrop;
use std::os::fd::FromRawFd;
let mut file = ManuallyDrop::new(unsafe { std::fs::File::from_raw_fd(STATISTICS_FD) });
file.write_all(bytes)
.and_then(|()| file.flush())
.map_err(|source| JobError::Statistics {
reason: format!("writing to descriptor {STATISTICS_FD}: {source}"),
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn values_accumulate_and_replace() {
let mut stats = JobStatistics::new();
stats.add("rows", 1).unwrap();
stats.add("rows", 2).unwrap();
assert_eq!(stats.get("rows"), Some(3));
stats.set("rows", 10).unwrap();
assert_eq!(stats.get("rows"), Some(10));
assert_eq!(stats.len(), 1);
}
#[test]
fn a_counter_that_runs_away_saturates_instead_of_panicking() {
let mut stats = JobStatistics::new();
stats.set("rows", i64::MAX).unwrap();
stats.add("rows", 1).unwrap();
assert_eq!(stats.get("rows"), Some(i64::MAX));
}
#[test]
fn the_hundred_and_twenty_ninth_name_is_refused() {
let mut stats = JobStatistics::new();
for i in 0..MAX_STATISTICS {
stats.add(format!("stat_{i}"), 1).unwrap();
}
stats.add("stat_0", 1).unwrap();
assert_eq!(stats.get("stat_0"), Some(2));
let err = stats.add("one_too_many", 1).expect_err("must refuse");
assert!(err.to_string().contains("128"), "{err}");
assert!(err.to_string().contains("one_too_many"), "{err}");
}
#[test]
fn the_encoding_is_a_yson_list_fragment_of_one_map() {
let mut stats = JobStatistics::new();
stats.set("rows/rejected", 7).unwrap();
stats.set("bytes", 4096).unwrap();
let encoded = String::from_utf8(stats.encode().unwrap()).unwrap();
assert_eq!(encoded, r#"{bytes=4096;"rows/rejected"=7};"#);
}
#[test]
fn nothing_is_encoded_for_no_statistics() {
let mut stats = JobStatistics::new();
assert!(stats.is_empty());
stats.finish().unwrap();
}
}