pitchfork-cli 2.26.0

Daemons with DX
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
use axum::{
    body::Body,
    extract::{Path, Query},
    http::StatusCode,
    response::{IntoResponse, Response},
};
use chrono::{DateTime, Local, TimeZone};
use serde::Deserialize;
use std::convert::Infallible;

use crate::cli::json_output::JsonLogEntry;
use crate::daemon_id::DaemonId;
use crate::log_store::sqlite::LOG_STORE;
use crate::log_store::{FieldFilter, LogQuery, LogStore, MessageFilter};

#[derive(Deserialize)]
pub struct TailQuery {
    lines: Option<usize>,
    since: Option<String>,
    until: Option<String>,
    level: Option<String>,
    grep: Option<String>,
    regex: Option<String>,
    logger: Option<String>,
    /// Structured field filters in "KEY=VALUE" format (can be repeated).
    field: Option<Vec<String>>,
    /// Whether grep should be case-sensitive. Default: false.
    case_sensitive: Option<bool>,
    /// jq expression for advanced filtering.
    jq: Option<String>,
    /// When set, return only entries with id < before_id (backward pagination
    /// for scroll-up history loading). The response is a one-shot JSONL body
    /// instead of a streaming response.
    before_id: Option<i64>,
}

/// Parse a datetime string from the query params.
/// Accepts ISO 8601, "YYYY-MM-DD HH:MM[:SS]", or "YYYY-MM-DDTHH:MM[:SS]".
/// The optional-seconds form matches HTML datetime-local input output.
fn parse_datetime(s: &str) -> Option<DateTime<Local>> {
    // Try ISO 8601 first
    if let Ok(dt) = DateTime::parse_from_rfc3339(s) {
        return Some(dt.with_timezone(&Local));
    }
    // Try "YYYY-MM-DD[T ]HH:MM[:SS]" — covers datetime-local (no seconds),
    // space-separated, and full datetime formats.
    for fmt in [
        "%Y-%m-%dT%H:%M:%S",
        "%Y-%m-%d %H:%M:%S",
        "%Y-%m-%dT%H:%M",
        "%Y-%m-%d %H:%M",
    ] {
        if let Ok(naive) = chrono::NaiveDateTime::parse_from_str(s, fmt) {
            return Local.from_local_datetime(&naive).single();
        }
    }
    None
}

/// Build message and field filters from query params.
/// Returns an error string if any filter value is invalid.
fn build_filters(query: &TailQuery) -> Result<(Vec<MessageFilter>, Vec<FieldFilter>), String> {
    let mut message_filters = Vec::new();
    let mut field_filters = Vec::new();

    if let Some(grep) = query.grep.as_deref().filter(|s| !s.is_empty()) {
        message_filters.push(MessageFilter::Contains {
            pattern: grep.to_string(),
            case_sensitive: query.case_sensitive.unwrap_or(false),
        });
    }

    if let Some(regex) = query.regex.as_deref().filter(|s| !s.is_empty()) {
        // Pre-validate regex so the user gets a clear 400 error instead of
        // a SQLite user-function failure at query time (which the polling
        // loop would silently swallow, stalling the stream).
        if let Err(e) = regex::Regex::new(regex) {
            return Err(format!("invalid regex pattern: {e}"));
        }
        message_filters.push(MessageFilter::Regex {
            pattern: regex.to_string(),
        });
    }

    if let Some(level) = query.level.as_deref().filter(|s| !s.is_empty()) {
        match crate::log_parse::normalize_level_str(level) {
            Some(normalized) => field_filters.push(FieldFilter::LevelMin(normalized)),
            None => return Err(format!("invalid level value: {level}")),
        }
    }

    if let Some(logger) = query.logger.as_deref().filter(|s| !s.is_empty()) {
        field_filters.push(FieldFilter::LoggerContains(logger.to_string()));
    }

    // Parse "KEY=VALUE" field filters (can be repeated).
    if let Some(fields) = &query.field {
        for pair in fields {
            let pair = pair.trim();
            if pair.is_empty() {
                continue;
            }
            let Some((key, value)) = pair.split_once('=') else {
                return Err(format!(
                    "invalid field filter: '{pair}' (expected KEY=VALUE format)"
                ));
            };
            if key.is_empty() {
                return Err("invalid field filter: empty key".to_string());
            }
            field_filters.push(FieldFilter::FieldEq {
                key: key.to_string(),
                value: value.to_string(),
            });
        }
    }

    Ok((message_filters, field_filters))
}

pub async fn tail(Path(id): Path<String>, Query(query): Query<TailQuery>) -> Response<Body> {
    let daemon_id = match DaemonId::parse(&id) {
        Ok(id) => id,
        Err(_) => {
            return Response::builder()
                .status(StatusCode::BAD_REQUEST)
                .header("content-type", "text/plain")
                .body(Body::from("invalid daemon id"))
                .unwrap();
        }
    };

    let history_lines = query.lines.unwrap_or(100);
    let qualified = daemon_id.qualified();

    // Parse time range. Invalid datetime strings return 400 instead of
    // silently broadening the query.
    let from = match query.since.as_deref().filter(|s| !s.is_empty()) {
        Some(s) => match parse_datetime(s) {
            Some(dt) => Some(dt),
            None => {
                return Response::builder()
                    .status(StatusCode::BAD_REQUEST)
                    .header("content-type", "text/plain")
                    .body(Body::from(format!("invalid 'since' datetime: {s}")))
                    .unwrap();
            }
        },
        None => None,
    };
    let to = match query.until.as_deref().filter(|s| !s.is_empty()) {
        Some(s) => match parse_datetime(s) {
            Some(dt) => Some(dt),
            None => {
                return Response::builder()
                    .status(StatusCode::BAD_REQUEST)
                    .header("content-type", "text/plain")
                    .body(Body::from(format!("invalid 'until' datetime: {s}")))
                    .unwrap();
            }
        },
        None => None,
    };
    let (message_filters, field_filters) = match build_filters(&query) {
        Ok(v) => v,
        Err(e) => {
            return Response::builder()
                .status(StatusCode::BAD_REQUEST)
                .header("content-type", "text/plain")
                .body(Body::from(e))
                .unwrap();
        }
    };

    // Compile jq filter early so parse errors surface before any query.
    let jq_filter = match query.jq.as_deref().filter(|s| !s.is_empty()) {
        Some(expr) => match crate::log_jq::JqFilter::new(expr) {
            Ok(f) => Some(f),
            Err(e) => {
                return Response::builder()
                    .status(StatusCode::BAD_REQUEST)
                    .header("content-type", "text/plain")
                    .body(Body::from(format!("invalid jq expression: {e}")))
                    .unwrap();
            }
        },
        None => None,
    };

    // Fetch initial history and clear generation atomically (single
    // connection lock + transaction) so a concurrent clear cannot pair
    // stale history with a new generation and evade clear detection.
    let (initial, initial_gen) = match tokio::task::spawn_blocking({
        let q = qualified.clone();
        let mf = message_filters.clone();
        let ff = field_filters.clone();
        let d = daemon_id.clone();
        move || {
            LOG_STORE.query_with_generation(
                &LogQuery {
                    daemon_ids: vec![q],
                    from,
                    to,
                    limit: Some(history_lines),
                    order_desc: true,
                    after_id: None,
                    before_id: query.before_id,
                    message_filters: mf,
                    field_filters: ff,
                    include_structured: true,
                },
                &d,
            )
        }
    })
    .await
    {
        Ok(Ok((entries, clear_gen))) => (entries, clear_gen),
        Ok(Err(e)) => {
            log::warn!("failed to query logs for {daemon_id}: {e}");
            return Response::builder()
                .status(StatusCode::NOT_FOUND)
                .header("content-type", "text/plain")
                .body(Body::from(format!("failed to query logs: {e}")))
                .unwrap();
        }
        Err(e) => {
            log::warn!("log query task panicked for {daemon_id}: {e}");
            return Response::builder()
                .status(StatusCode::INTERNAL_SERVER_ERROR)
                .header("content-type", "text/plain")
                .body(Body::from("log query failed"))
                .unwrap();
        }
    };

    // When before_id is set, this is a one-shot backward pagination request
    // (scroll-up history loading). Return the results immediately without
    // starting the streaming polling loop.
    if query.before_id.is_some() {
        let initial = if let Some(jq) = &jq_filter {
            jq.filter(initial)
        } else {
            initial
        };
        let body: String = initial
            .into_iter()
            .rev()
            .filter_map(|e| {
                let entry: JsonLogEntry = e.into();
                serde_json::to_string(&entry).ok().map(|s| s + "\n")
            })
            .collect();
        return Response::builder()
            .status(StatusCode::OK)
            .header("content-type", "application/x-ndjson")
            .header("x-log-generation", initial_gen.unwrap_or(0).to_string())
            .body(Body::from(body))
            .unwrap();
    }

    // Capture cursor from raw query (before jq filtering) so the polling
    // loop doesn't rescan jq-filtered-out rows on every poll.
    let cursor_id = initial.first().map(|e| e.id).unwrap_or(0);

    // Apply jq filter if present.
    let initial = if let Some(jq) = &jq_filter {
        jq.filter(initial)
    } else {
        initial
    };

    // Reverse so oldest lines are yielded first.
    let initial: Vec<String> = initial
        .into_iter()
        .rev()
        .filter_map(|e| {
            let entry: JsonLogEntry = e.into();
            serde_json::to_string(&entry).ok().map(|s| s + "\n")
        })
        .collect();

    let qualified_clone = qualified.clone();
    let stream = async_stream::stream! {
        // Yield history
        for line in initial {
            yield Ok::<Vec<u8>, Infallible>(line.into_bytes());
        }

        let mut last_id: i64 = cursor_id;

        // initial_gen was captured atomically with the initial history query.
        // A value of None means the daemon has never been cleared (no row in
        // log_clear_generations). Treat that as generation 0 so a subsequent
        // clear (which bumps to 1+) is detected as a change.
        let mut last_clear_gen: u64 = initial_gen.unwrap_or(0);

        const BATCH_SIZE: usize = 500;
        loop {
            tokio::time::sleep(tokio::time::Duration::from_millis(500)).await;

            // Read new rows and current generation atomically (single
            // connection lock + transaction) so a clear cannot interleave
            // between the generation check and the row query. Without this,
            // a clear could happen after the generation read but before the
            // row query, producing post-clear rows that later get replayed
            // when the next poll detects the generation change.
            let poll_result = match tokio::task::spawn_blocking({
                let q = qualified_clone.clone();
                let mf = message_filters.clone();
                let ff = field_filters.clone();
                let d = daemon_id.clone();
                move || LOG_STORE.query_with_generation(
                    &LogQuery {
                        daemon_ids: vec![q],
                        from,
                        to,
                        limit: Some(BATCH_SIZE),
                        order_desc: false,
                        after_id: Some(last_id),
                        before_id: None,
                        message_filters: mf,
                        field_filters: ff,
                        include_structured: true,
                    },
                    &d,
                )
            })
            .await
            {
                Ok(Ok((entries, generation))) => (entries, generation.unwrap_or(0)),
                _ => continue,
            };

            let (raw_entries, current_gen) = poll_result;

            if current_gen != last_clear_gen {
                // Log clear detected — signal the frontend to flush its
                // buffer, then reset cursor and skip this batch.
                last_clear_gen = current_gen;
                last_id = 0;
                yield Ok::<Vec<u8>, Infallible>(
                    format!("{{\"_clear\":true,\"_gen\":{}}}\n", current_gen).into_bytes(),
                );
                continue;
            }

            // Advance cursor past all raw entries (not just jq-matched ones)
            // so filtered-out rows aren't re-scanned on every poll.
            if let Some(last) = raw_entries.last() {
                last_id = last.id;
            }

            // Apply jq filter if present.
            let entries = if let Some(jq) = &jq_filter {
                jq.filter(raw_entries)
            } else {
                raw_entries
            };

            for entry in entries {
                let json_entry: JsonLogEntry = entry.into();
                if let Ok(line) = serde_json::to_string(&json_entry) {
                    yield Ok::<Vec<u8>, Infallible>((line + "\n").into_bytes());
                }
            }
        }
    };

    Response::builder()
        .status(StatusCode::OK)
        .header("content-type", "application/x-ndjson")
        .header("x-log-generation", initial_gen.unwrap_or(0).to_string())
        .body(Body::from_stream(stream))
        .unwrap()
}

/// Return distinct logger names for a daemon, for populating filter dropdowns.
pub async fn loggers(Path(id): Path<String>) -> Response<Body> {
    let daemon_id = match DaemonId::parse(&id) {
        Ok(id) => id,
        Err(_) => {
            return Response::builder()
                .status(StatusCode::BAD_REQUEST)
                .header("content-type", "text/plain")
                .body(Body::from("invalid daemon id"))
                .unwrap();
        }
    };

    let qualified = daemon_id.qualified();
    let loggers =
        match tokio::task::spawn_blocking(move || LOG_STORE.distinct_loggers(&qualified)).await {
            Ok(Ok(loggers)) => loggers,
            Ok(Err(e)) => {
                log::warn!("failed to query loggers for {daemon_id}: {e}");
                return Response::builder()
                    .status(StatusCode::INTERNAL_SERVER_ERROR)
                    .header("content-type", "text/plain")
                    .body(Body::from("failed to query loggers"))
                    .unwrap();
            }
            Err(e) => {
                log::warn!("loggers query task panicked for {daemon_id}: {e}");
                return Response::builder()
                    .status(StatusCode::INTERNAL_SERVER_ERROR)
                    .header("content-type", "text/plain")
                    .body(Body::from("loggers query failed"))
                    .unwrap();
            }
        };

    (
        StatusCode::OK,
        [("content-type", "application/json")],
        serde_json::to_string(&loggers).unwrap_or_else(|_| "[]".to_string()),
    )
        .into_response()
}

/// Return distinct structured field keys for a daemon, for jq autocomplete.
pub async fn field_keys(Path(id): Path<String>) -> Response<Body> {
    let daemon_id = match DaemonId::parse(&id) {
        Ok(id) => id,
        Err(_) => {
            return Response::builder()
                .status(StatusCode::BAD_REQUEST)
                .header("content-type", "text/plain")
                .body(Body::from("invalid daemon id"))
                .unwrap();
        }
    };

    let qualified = daemon_id.qualified();
    let keys = match tokio::task::spawn_blocking(move || LOG_STORE.distinct_field_keys(&qualified))
        .await
    {
        Ok(Ok(keys)) => keys,
        Ok(Err(e)) => {
            log::warn!("failed to query field keys for {daemon_id}: {e}");
            return Response::builder()
                .status(StatusCode::INTERNAL_SERVER_ERROR)
                .header("content-type", "text/plain")
                .body(Body::from("failed to query field keys"))
                .unwrap();
        }
        Err(e) => {
            log::warn!("field keys query task panicked for {daemon_id}: {e}");
            return Response::builder()
                .status(StatusCode::INTERNAL_SERVER_ERROR)
                .header("content-type", "text/plain")
                .body(Body::from("field keys query failed"))
                .unwrap();
        }
    };

    (
        StatusCode::OK,
        [("content-type", "application/json")],
        serde_json::to_string(&keys).unwrap_or_else(|_| "[]".to_string()),
    )
        .into_response()
}