Skip to main content

lava_api/
joblog.rs

1use std::collections::HashMap;
2use std::fmt;
3use std::pin::Pin;
4use std::task::Poll;
5use std::time::Duration;
6
7use bytes::{Bytes, BytesMut};
8use chrono::{DateTime, NaiveDateTime, Utc};
9use futures::future::BoxFuture;
10use futures::stream::BoxStream;
11use futures::{prelude::*, ready};
12use reqwest::{Response, StatusCode, Url};
13use serde::{Deserialize, Deserializer};
14use thiserror::Error;
15
16use crate::Lava;
17
18#[derive(Debug)]
19pub struct JobLogBuilder<'a> {
20    lava: &'a Lava,
21    id: i64,
22    start: u64,
23    end: u64,
24}
25
26impl<'a> JobLogBuilder<'a> {
27    pub fn new(lava: &'a Lava, id: i64) -> Self {
28        Self {
29            lava,
30            id,
31            start: 0,
32            end: 0,
33        }
34    }
35
36    pub fn start(mut self, start: u64) -> Self {
37        self.start = start;
38        self
39    }
40
41    pub fn end(mut self, end: u64) -> Self {
42        self.end = end;
43        self
44    }
45
46    pub fn raw(self) -> JobLogRaw<'a> {
47        JobLogRaw::new(self.lava, self.id, self.start, self.end)
48    }
49
50    pub fn log(self) -> JobLog<'a> {
51        JobLog::new(self.lava, self.id, self.start, self.end)
52    }
53}
54
55#[derive(Debug, Error)]
56pub enum JobLogError {
57    #[error("Request failed: {0}")]
58    RequestError(#[from] reqwest::Error),
59    #[error("Parse error: {0} - {1}")]
60    ParseError(String, serde_norway::Error),
61    #[error("No data available")]
62    NoData,
63}
64
65enum LogRequest {
66    Initial,
67    Request(BoxFuture<'static, reqwest::Result<Response>>),
68    Stream(BoxStream<'static, reqwest::Result<Bytes>>),
69    Done,
70}
71
72impl fmt::Debug for LogRequest {
73    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
74        let fmt = match self {
75            LogRequest::Initial => "Initial",
76            LogRequest::Request(_) => "Request",
77            LogRequest::Stream(_) => "Stream",
78            LogRequest::Done => "Done",
79        };
80        f.write_str(fmt)
81    }
82}
83
84#[derive(Debug)]
85pub struct JobLogRaw<'a> {
86    lava: &'a Lava,
87    id: i64,
88    start: u64,
89    end: u64,
90    state: LogRequest,
91}
92
93impl<'a> JobLogRaw<'a> {
94    fn new(lava: &'a Lava, id: i64, start: u64, end: u64) -> Self {
95        Self {
96            lava,
97            id,
98            start,
99            end,
100            state: LogRequest::Initial,
101        }
102    }
103
104    fn url(&self) -> Url {
105        let mut url = self.lava.base.clone();
106        url.path_segments_mut()
107            .unwrap()
108            .pop_if_empty()
109            .push("jobs")
110            .push(&self.id.to_string())
111            .push("logs")
112            .push("");
113
114        if self.start != 0 {
115            url.query_pairs_mut()
116                .append_pair("start", &self.start.to_string());
117        }
118
119        if self.end != 0 {
120            url.query_pairs_mut()
121                .append_pair("end", &self.end.to_string());
122        }
123        url
124    }
125}
126
127impl Stream for JobLogRaw<'_> {
128    type Item = Result<Bytes, JobLogError>;
129
130    fn poll_next(
131        self: std::pin::Pin<&mut Self>,
132        cx: &mut std::task::Context<'_>,
133    ) -> std::task::Poll<Option<Self::Item>> {
134        let me = self.get_mut();
135        loop {
136            match me.state {
137                LogRequest::Initial => {
138                    let u = me.url();
139                    let r = me.lava.client.get(u).send();
140                    me.state = LogRequest::Request(r.boxed());
141                }
142                LogRequest::Request(ref mut r) => match ready!(r.as_mut().poll(cx)) {
143                    Ok(r) => match r.error_for_status() {
144                        Ok(r) => me.state = LogRequest::Stream(r.bytes_stream().boxed()),
145                        Err(e) => {
146                            me.state = LogRequest::Done;
147                            let e = match e.status() {
148                                Some(StatusCode::NOT_FOUND) => JobLogError::NoData,
149                                _ => e.into(),
150                            };
151                            return Poll::Ready(Some(Err(e)));
152                        }
153                    },
154                    Err(e) => return Poll::Ready(Some(Err(e.into()))),
155                },
156                LogRequest::Stream(ref mut stream) => match ready!(stream.as_mut().poll_next(cx)) {
157                    Some(Err(e)) => return Poll::Ready(Some(Err(e.into()))),
158                    Some(Ok(b)) => {
159                        return Poll::Ready(Some(Ok(b)));
160                    }
161                    None => {
162                        me.state = LogRequest::Done;
163                        return Poll::Ready(None);
164                    }
165                },
166                LogRequest::Done => return Poll::Ready(None),
167            }
168        }
169    }
170
171    fn size_hint(&self) -> (usize, Option<usize>) {
172        (0, None)
173    }
174}
175
176fn deserialize_duration<'de, D>(d: D) -> Result<Option<Duration>, D::Error>
177where
178    D: Deserializer<'de>,
179{
180    let duration = String::deserialize(d)?
181        .parse()
182        .map_err(serde::de::Error::custom)?;
183    Ok(Some(Duration::from_secs_f64(duration)))
184}
185
186#[derive(Debug, Clone, Deserialize)]
187pub struct JobResult {
188    pub case: String,
189    pub definition: String,
190    pub namespace: Option<String>,
191    pub level: Option<String>,
192    pub result: String,
193    #[serde(default, deserialize_with = "deserialize_duration")]
194    pub duration: Option<Duration>,
195    #[serde(default)]
196    pub extra: HashMap<String, serde_norway::Value>,
197}
198
199#[derive(Debug, Clone)]
200pub enum JobLogMsg {
201    Msg(String),
202    Msgs(Vec<String>),
203    Result(JobResult),
204}
205
206impl<'de> Deserialize<'de> for JobLogMsg {
207    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
208    where
209        D: Deserializer<'de>,
210    {
211        // Deserialize as serde_norway::Value first to avoid serde's ContentVisitor,
212        // which does not support YAML tagged values (e.g. `! "188250"`) that LAVA
213        // sometimes emits in the `extra` field.
214        let value = serde_norway::Value::deserialize(deserializer)?;
215        match value {
216            serde_norway::Value::String(s) => Ok(JobLogMsg::Msg(s)),
217            serde_norway::Value::Sequence(seq) => {
218                let msgs = seq
219                    .into_iter()
220                    .map(|v| match v {
221                        serde_norway::Value::String(s) => Ok(s),
222                        _ => Err(serde::de::Error::custom("expected string in msg sequence")),
223                    })
224                    .collect::<Result<Vec<_>, _>>()?;
225                Ok(JobLogMsg::Msgs(msgs))
226            }
227            serde_norway::Value::Mapping(_) => {
228                let result = serde_norway::from_value(value).map_err(serde::de::Error::custom)?;
229                Ok(JobLogMsg::Result(result))
230            }
231            _ => Err(serde::de::Error::custom(
232                "expected string, sequence, or mapping for msg",
233            )),
234        }
235    }
236}
237
238#[derive(Debug, Clone, Deserialize)]
239#[serde(rename_all = "lowercase")]
240pub enum JobLogLevel {
241    Debug,
242    Info,
243    Warning,
244    Error,
245    Results,
246    Target,
247    Input,
248    Feedback,
249    Exception,
250}
251
252fn job_deserialize_dt<'de, D>(d: D) -> Result<NaiveDateTime, D::Error>
253where
254    D: Deserializer<'de>,
255{
256    #[derive(Deserialize)]
257    #[serde(untagged)]
258    enum DorN {
259        D(DateTime<Utc>),
260        N(NaiveDateTime),
261    }
262    match DorN::deserialize(d)? {
263        DorN::D(d) => Ok(d.naive_utc()),
264        DorN::N(n) => Ok(n),
265    }
266}
267
268#[derive(Debug, Clone, Deserialize)]
269pub struct JobLogEntry {
270    #[serde(deserialize_with = "job_deserialize_dt")]
271    pub dt: NaiveDateTime,
272    pub lvl: JobLogLevel,
273    pub ns: Option<String>,
274    pub msg: JobLogMsg,
275}
276
277#[derive(Debug)]
278pub struct JobLog<'a> {
279    buf: Vec<Bytes>,
280    from_buf: bool,
281    raw: JobLogRaw<'a>,
282}
283
284impl<'a> JobLog<'a> {
285    fn new(lava: &'a Lava, id: i64, start: u64, end: u64) -> Self {
286        let raw = JobLogRaw::new(lava, id, start, end);
287        Self {
288            buf: Vec::new(),
289            from_buf: false,
290            raw,
291        }
292    }
293}
294
295impl Stream for JobLog<'_> {
296    type Item = Result<JobLogEntry, JobLogError>;
297
298    fn poll_next(
299        self: std::pin::Pin<&mut Self>,
300        cx: &mut std::task::Context<'_>,
301    ) -> Poll<Option<Self::Item>> {
302        let me = self.get_mut();
303        loop {
304            if me.from_buf {
305                let last = me.buf.last().unwrap();
306                if let Some(eol) = last.iter().position(|e| e == &b'\n') {
307                    let line = if me.buf.len() == 1 {
308                        if last.len() - 1 == eol {
309                            me.from_buf = false;
310                            me.buf.pop().unwrap()
311                        } else {
312                            let b = me.buf.get_mut(0).unwrap();
313                            b.split_to(eol + 1)
314                        }
315                    } else {
316                        let mut buf = BytesMut::new();
317                        for b in me.buf.drain(0..me.buf.len() - 1) {
318                            buf.extend_from_slice(b.as_ref());
319                        }
320
321                        let last = me.buf.last().unwrap();
322                        if last.len() == eol {
323                            me.from_buf = false;
324                            buf.extend_from_slice(me.buf.pop().unwrap().as_ref());
325                        } else {
326                            let b = me.buf.get_mut(0).unwrap();
327                            buf.extend_from_slice(b.split_to(eol + 1).as_ref());
328                        }
329                        buf.into()
330                    };
331                    let l = line.slice(1..);
332                    let entry = serde_norway::from_slice(l.as_ref()).map_err(|e| {
333                        let s = String::from_utf8_lossy(l.as_ref());
334                        JobLogError::ParseError(s.into_owned(), e)
335                    });
336                    return Poll::Ready(Some(entry));
337                } else {
338                    me.from_buf = false;
339                }
340            } else {
341                match ready!(Pin::new(&mut me.raw).poll_next(cx)) {
342                    Some(Err(e)) => return Poll::Ready(Some(Err(e))),
343                    Some(Ok(b)) => {
344                        me.from_buf = true;
345                        me.buf.push(b);
346                    }
347                    None => return Poll::Ready(None),
348                }
349            }
350        }
351    }
352}
353
354#[cfg(test)]
355mod tests {
356    use super::*;
357
358    #[test]
359    fn job_log_entry_deserialisation() {
360        let entry0: JobLogEntry = serde_json::from_str(
361            r#"{"dt": "2026-06-02T12:16:27.730458", "lvl": "results", "msg": {"definition": "lava", "case": "job", "result": "pass"}}"#,
362        )
363        .unwrap();
364        let entry1: JobLogEntry = serde_json::from_str(
365            r#"{"dt": "2026-06-02T12:16:27.730458+00:00", "lvl": "results", "msg": {"definition": "lava", "case": "job", "result": "pass"}}"#,
366        )
367        .unwrap();
368        assert_eq!(entry0.dt, entry1.dt);
369    }
370
371    #[test]
372    fn test_joblog_entry_with_tagged_extra() {
373        // Regression test for https://github.com/collabora/lava-api/issues/31
374        // LAVA sometimes emits tagged YAML values (e.g. `! "188250"`) in the
375        // `extra` field of a results entry. The bare `!` is a YAML non-specific
376        // tag and should not prevent deserialization.
377        let yaml = r#"{"dt": "2026-04-30T16:19:53.641343", "lvl": "results", "msg": {"definition": "lava", "namespace": "common", "case": "http-download", "level": "1.4.1", "duration": "0.00", "result": "pass", "extra": {"label": "dtb", "size": ! "188250", "sha256sum": "9f4c38d2218c38be09a03812ab9176d521b8fb624b8e9ecff69b4125d260122b"}}}"#;
378        let entry: JobLogEntry =
379            serde_norway::from_str(yaml).expect("failed to parse entry with tagged extra value");
380        assert!(matches!(entry.lvl, JobLogLevel::Results));
381        let JobLogMsg::Result(result) = entry.msg else {
382            panic!("expected Result variant, got {:?}", entry.msg);
383        };
384        assert_eq!(result.case, "http-download");
385        assert_eq!(result.definition, "lava");
386        assert_eq!(result.result, "pass");
387        assert!(result.extra.contains_key("size"));
388        assert!(result.extra.contains_key("sha256sum"));
389    }
390
391    #[test]
392    fn test_joblog_msg_string() {
393        let yaml = r#"{"dt": "2022-01-01T00:00:00", "lvl": "info", "msg": "hello world"}"#;
394        let entry: JobLogEntry = serde_norway::from_str(yaml).expect("failed to parse string msg");
395        assert!(matches!(entry.msg, JobLogMsg::Msg(ref s) if s == "hello world"));
396    }
397
398    #[test]
399    fn test_joblog_msg_sequence() {
400        let yaml = r#"{"dt": "2022-01-01T00:00:00", "lvl": "target", "msg": ["line1", "line2"]}"#;
401        let entry: JobLogEntry =
402            serde_norway::from_str(yaml).expect("failed to parse sequence msg");
403        assert!(matches!(entry.msg, JobLogMsg::Msgs(ref v) if v == &["line1", "line2"]));
404    }
405}