drain_flow/drains/
simple.rs

1// Copyright Nicholas Harring. All rights reserved.
2//
3// This program is free software: you can redistribute it and/or modify it under
4// the terms of the Server Side Public License, version 1, as published by MongoDB, Inc.
5// This program is distributed in the hope that it will be useful, but WITHOUT ANY WARRANTY;
6// without even the implied warranty of MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.
7// See the Server Side Public License for more details. You should have received a copy of the
8// Server Side Public License along with this program.
9// If not, see <http://www.mongodb.com/licensing/server-side-public-license>.
10
11use std::{collections::HashMap, fmt, sync::Arc};
12
13use anyhow::{anyhow, Error};
14use fraction::{BigInt, FromPrimitive, Ratio};
15use joinery::{Joinable, JoinableIterator};
16use lazy_static::lazy_static;
17use parking_lot::RwLock;
18use regex::Regex;
19use string_interner::{DefaultSymbol, StringInterner};
20use tracing::instrument;
21
22use crate::{drains::api::Drain, log_group::LogGroup, record::Record};
23
24lazy_static! {
25    pub(crate) static ref INTERNER: Arc<RwLock<StringInterner<string_interner::backend::BucketBackend>>> =
26        Arc::new(RwLock::new(StringInterner::<
27            string_interner::backend::BucketBackend,
28        >::new()));
29}
30#[derive(Debug, Clone)]
31pub struct SingleLayer {
32    pub domain: Vec<Regex>,
33    // NumTokens -> First Token -> List of Log groups
34    base_layer: HashMap<usize, HashMap<DefaultSymbol, Vec<LogGroup>>>,
35    pub threshold: Ratio<BigInt>,
36    strings: Arc<RwLock<StringInterner<string_interner::backend::BucketBackend>>>,
37}
38
39impl SingleLayer {
40    #[instrument(skip(domain))]
41    pub fn new(domain: Vec<String>) -> Result<Self, Error> {
42        let patterns = domain
43            .iter()
44            .map(|s| Regex::new(s))
45            .collect::<Result<Vec<Regex>, regex::Error>>()?;
46        Ok(Self {
47            domain: patterns,
48            base_layer: HashMap::new(),
49            threshold: Ratio::from_float::<f32>(0.5).expect("0.5 converts into a ratio"),
50            strings: INTERNER.clone(),
51        })
52    }
53
54    #[instrument(skip(self))]
55    pub fn set_threshold(&mut self, numerator: u64, denominator: u64) -> Result<(), Error> {
56        let numer = BigInt::from_u64(numerator)
57            .ok_or_else(|| anyhow!("unable to make numerator from {}", numerator))?;
58        let denom = BigInt::from_u64(denominator)
59            .ok_or_else(|| anyhow!("unable to make denominator from {}", denominator))?;
60        let new_ratio = Ratio::new(numer, denom);
61        self.threshold = new_ratio;
62        Ok(())
63    }
64
65    #[instrument(skip(self), level = "trace")]
66    fn iter_groups(&self) -> Vec<Vec<&LogGroup>> {
67        let mut results: Vec<Vec<&LogGroup>> = Vec::new();
68        for length in self.base_layer.keys() {
69            let mut groups = vec![];
70            for (_, grp) in self.base_layer.get(length).unwrap().iter() {
71                for g in grp {
72                    groups.push(g);
73                }
74            }
75            results.push(groups);
76        }
77        results
78    }
79
80    #[instrument(skip(self), level = "trace")]
81    pub fn resolve(&self, sym: DefaultSymbol) -> String {
82        self.strings
83            .read()
84            .resolve(sym)
85            .expect("symbols must resolve")
86            .to_owned()
87    }
88}
89
90impl Drain for SingleLayer {
91    /// Accepts a line of input for processing against existing records
92    ///
93    /// Return
94    /// Ok(true) when a new entry is added
95    /// Ok(false) when the line matched an existing entry
96    /// Err(e) for errors during processing
97    #[instrument(skip(self, line))]
98    fn process_line(&mut self, line: String) -> Result<bool, Error> {
99        if line.is_empty() {
100            return Ok(false);
101        }
102        let new_record = Record::new(line);
103        let length = new_record.len();
104        let first = new_record.first().expect("records have first tokens");
105        if let Some(second_layer) = self.base_layer.get_mut(&length) {
106            match second_layer.get_mut(&first) {
107                Some(log_groups) => {
108                    let (score, offset) = log_groups.iter_mut().enumerate().fold(
109                        (
110                            0, // best score
111                            0, // index of best score LogGroup
112                        ),
113                        |mut acc, elem| {
114                            let score = new_record.clone().calc_sim_score(elem.1.event());
115                            if score > acc.0 {
116                                acc = (score, elem.0); // overwrite state with new values
117                            }
118                            acc
119                        },
120                    );
121                    let score_ratio =
122                        Ratio::<BigInt>::new(BigInt::from(score), BigInt::from(length));
123                    if score_ratio > self.threshold {
124                        // add this record's uid to the list of examples for the log group
125                        log_groups[offset].add_example(new_record);
126                        Ok(false)
127                    } else {
128                        log_groups.push(LogGroup::new(new_record));
129                        Ok(true)
130                    }
131                }
132                None => {
133                    second_layer.insert(first, vec![LogGroup::new(new_record)]);
134                    Ok(true)
135                }
136            }
137        } else {
138            self.base_layer.insert(length, HashMap::new());
139            let second_layer = self
140                .base_layer
141                .get_mut(&length)
142                .expect("We just inserted this map");
143            second_layer.insert(first, vec![LogGroup::new(new_record)]);
144            Ok(true)
145        }
146    }
147
148    fn collect_log_groups(&self) -> Vec<LogGroup> {
149        self.iter_groups().into_iter().flatten().cloned().collect()
150    }
151}
152
153impl fmt::Display for SingleLayer {
154    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
155        let base = format!(
156            "SimpleDrain\nDomain Patterns: {:?}\nSimilarity Threshold: {}\n",
157            self.domain, self.threshold
158        );
159        let lg = "Log Groups:\n".to_string();
160        let groups = self
161            .iter_groups()
162            .iter()
163            .flatten()
164            .map(std::string::ToString::to_string)
165            .collect::<Vec<String>>();
166        let group_str = groups.iter().join_with("\n");
167        write!(f, "{}", [base, lg, group_str.to_string()].join_concat())
168    }
169}
170
171#[cfg(test)]
172mod should {
173    use spectral::prelude::*;
174    use tracing_test::traced_test;
175
176    use crate::drains::{api::Drain, simple::SingleLayer}; // Removed <BucketBackend> for now, will add if compiler complains
177
178    #[traced_test]
179    #[test]
180    fn test_new_drain() {
181        let drain = SingleLayer::new(vec![]);
182        assert_that(&drain).is_ok();
183    }
184
185    #[traced_test]
186    #[test]
187    fn test_set_threshold() {
188        let mut drain = SingleLayer::new(vec![]).unwrap();
189        let res = drain.set_threshold(100, 200);
190        assert_that(&res).is_ok();
191    }
192
193    #[traced_test]
194    #[test]
195    fn test_single_process_line() {
196        let mut drain = SingleLayer::new(vec![]).unwrap();
197        let line_1 = "Message send failed to remote host: foo.bar.com".to_string();
198        let res = drain.process_line(line_1);
199        assert_that(&res).is_ok_containing(true);
200    }
201
202    #[traced_test]
203    #[test]
204    fn test_multiple_process_line() {
205        let mut drain = SingleLayer::new(vec![]).unwrap();
206        let line_1 = "Message send failed to remote host: foo.bar.com".to_string();
207        let line_2 = "Message send failed to remote host: bork.bork.com".to_string();
208        let line_3 = "Unknown error received from peer".to_string();
209        let res = drain.process_line(line_1);
210        assert_that(&res).is_ok_containing(true);
211        let res = drain.process_line(line_2);
212        assert_that(&res).is_ok_containing(false);
213        let res = drain.process_line(line_3);
214        assert_that!(res).is_ok_containing(true);
215    }
216
217    #[traced_test]
218    #[test]
219    fn test_iter_groups() {
220        let line_1 = "This is a sequence".to_string();
221        let line_2 = "Another different order of words".to_string();
222        let line_3 = "Finally one last unique set of character runs".to_string();
223        let mut drain = SingleLayer::new(vec![]).unwrap();
224        Drain::process_line(&mut drain, line_1).unwrap();
225        Drain::process_line(&mut drain, line_2).unwrap();
226        Drain::process_line(&mut drain, line_3).unwrap();
227        let groups = drain.collect_log_groups();
228        assert_that(&groups).has_length(3);
229    }
230}