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 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 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}