lean_ctx/core/patterns/
spark.rs1use crate::core::compressor::strip_ansi;
12use std::collections::HashSet;
13
14pub fn compress(_cmd: &str, output: &str) -> Option<String> {
15 let trimmed = output.trim();
16 if trimmed.is_empty() {
17 return Some("spark: ok".to_string());
18 }
19
20 let mut jobs: Vec<String> = Vec::new();
21 let mut warnings: Vec<String> = Vec::new();
22 let mut warn_seen: HashSet<String> = HashSet::new();
23 let mut errors: Vec<String> = Vec::new();
24 let mut app_output: Vec<String> = Vec::new();
25 let mut saw_log = false;
26
27 for raw in trimmed.lines() {
28 let stripped = strip_ansi(raw);
29 let line = stripped.trim();
30 if line.is_empty() {
31 continue;
32 }
33
34 if let Some((level, rest)) = parse_log(line) {
35 saw_log = true;
36 match level {
37 "INFO" => {
38 if rest.contains("Job ") && rest.contains("finished") {
39 jobs.push(rest.to_string());
40 }
41 }
42 "WARN" => {
43 if warn_seen.insert(normalize(rest)) {
44 warnings.push(rest.to_string());
45 }
46 }
47 "ERROR" => errors.push(rest.to_string()),
48 _ => {}
49 }
50 } else if is_exception(line) {
51 errors.push(line.to_string());
52 } else {
53 app_output.push(line.to_string());
55 }
56 }
57
58 if !saw_log && jobs.is_empty() && errors.is_empty() {
59 return Some(fallback(trimmed));
60 }
61
62 let mut parts = vec![format!(
63 "spark: {} job(s), {} warning(s), {} error(s)",
64 jobs.len(),
65 warnings.len(),
66 errors.len()
67 )];
68 push_capped(&mut parts, &jobs, 10, "more jobs");
69 push_capped(&mut parts, &warnings, 5, "more warnings");
70 push_capped(&mut parts, &errors, 10, "more errors");
71 push_capped(&mut parts, &app_output, 20, "more output lines");
72 Some(parts.join("\n"))
73}
74
75fn parse_log(line: &str) -> Option<(&str, &str)> {
77 let mut it = line.splitn(4, ' ');
78 let date = it.next()?;
79 let time = it.next()?;
80 let level = it.next()?;
81 let rest = it.next().unwrap_or("");
82 if looks_like_date(date) && looks_like_time(time) && is_level(level) {
83 Some((level, rest))
84 } else {
85 None
86 }
87}
88
89fn looks_like_date(s: &str) -> bool {
90 let parts: Vec<&str> = s.split('/').collect();
91 parts.len() == 3 && parts.iter().all(|p| p.chars().all(|c| c.is_ascii_digit()))
92}
93
94fn looks_like_time(s: &str) -> bool {
95 let parts: Vec<&str> = s.split(':').collect();
96 parts.len() == 3 && parts.iter().all(|p| p.chars().all(|c| c.is_ascii_digit()))
97}
98
99fn is_level(s: &str) -> bool {
100 matches!(s, "INFO" | "WARN" | "ERROR" | "DEBUG" | "TRACE")
101}
102
103fn is_exception(line: &str) -> bool {
104 line.contains("Exception") || line.starts_with("Caused by:") || line.contains("Error:")
105}
106
107fn normalize(s: &str) -> String {
109 s.chars().filter(|c| !c.is_ascii_digit()).collect()
110}
111
112fn push_capped(parts: &mut Vec<String>, items: &[String], cap: usize, label: &str) {
113 for item in items.iter().take(cap) {
114 parts.push(format!(" {item}"));
115 }
116 if items.len() > cap {
117 parts.push(format!(" ... +{} {label}", items.len() - cap));
118 }
119}
120
121fn fallback(text: &str) -> String {
122 let lines: Vec<&str> = text.lines().filter(|l| !l.trim().is_empty()).collect();
123 let n = lines.len().min(8);
124 let mut s = lines[..n].join("\n");
125 if lines.len() > n {
126 s.push_str(&format!("\n... (+{} lines)", lines.len() - n));
127 }
128 s
129}
130
131#[cfg(test)]
132mod tests {
133 use super::*;
134
135 const LOG: &str = "23/01/01 12:00:00 INFO SparkContext: Running Spark version 3.4.0\n23/01/01 12:00:01 INFO ResourceUtils: No custom resources configured\n23/01/01 12:00:02 INFO Utils: Successfully started service\n23/01/01 12:00:03 WARN NativeCodeLoader: Unable to load native-hadoop\n23/01/01 12:00:10 INFO DAGScheduler: Job 0 finished: collect at Main.scala:20, took 5.123 s\n23/01/01 12:00:15 ERROR Executor: Exception in task 0.0 in stage 1.0\n";
136
137 #[test]
138 fn drops_info_keeps_job_warn_error() {
139 let r = compress("spark-submit app.py", LOG).unwrap();
140 assert!(r.contains("1 job(s), 1 warning(s), 1 error(s)"), "{r}");
141 assert!(r.contains("Job 0 finished"), "{r}");
142 assert!(r.contains("Executor: Exception"), "{r}");
143 assert!(!r.contains("ResourceUtils"), "drops info noise: {r}");
144 assert!(!r.contains("23/01/01"), "drops timestamps: {r}");
145 }
146
147 #[test]
148 fn keeps_application_stdout() {
149 let log = "23/01/01 12:00:00 INFO SparkContext: Running Spark version 3.4.0\n23/01/01 12:00:01 INFO ResourceUtils: No custom resources\n23/01/01 12:00:10 INFO DAGScheduler: Job 0 finished: collect, took 5.1 s\nResult: total words = 184273\n23/01/01 12:00:11 INFO SparkContext: Successfully stopped";
150 let r = compress("spark-submit app.py", log).unwrap();
151 assert!(
152 r.contains("Result: total words = 184273"),
153 "keeps app output: {r}"
154 );
155 assert!(!r.contains("ResourceUtils"), "still drops info noise: {r}");
156 }
157
158 #[test]
159 fn shorter_than_input() {
160 let r = compress("spark-submit app.py", LOG).unwrap();
161 assert!(r.len() < LOG.len());
162 }
163
164 #[test]
165 fn empty_is_ok() {
166 assert_eq!(compress("spark-submit app.py", "").unwrap(), "spark: ok");
167 }
168}