diff --git a/src/commands/logs.rs b/src/commands/logs.rs index 7a32b5f4..66228749 100644 --- a/src/commands/logs.rs +++ b/src/commands/logs.rs @@ -33,6 +33,8 @@ pub struct SearchArgs { pub from: String, pub to: String, pub limit: i32, + pub cursor: Option, + pub pages: u32, pub sort: String, pub storage: Option, pub index: Vec, @@ -204,16 +206,49 @@ fn parse_logs_sort(sort: &str) -> LogsSort { } } +fn append_log_page(aggregated: &mut Option, page: serde_json::Value) { + if let Some(aggregated) = aggregated { + if let Some(page_data) = page.get("data").and_then(|value| value.as_array()) { + if let Some(data) = aggregated + .get_mut("data") + .and_then(|value| value.as_array_mut()) + { + data.extend(page_data.iter().cloned()); + } else { + aggregated["data"] = serde_json::Value::Array(page_data.to_vec()); + } + } + if let Some(meta) = page.get("meta") { + aggregated["meta"] = meta.clone(); + } + } else { + *aggregated = Some(page); + } +} + +fn next_log_cursor(response: &serde_json::Value) -> Option { + response + .pointer("/meta/page/after") + .and_then(|value| value.as_str()) + .filter(|cursor| !cursor.is_empty()) + .map(str::to_owned) +} + pub async fn search(cfg: &Config, args: SearchArgs) -> Result<()> { let SearchArgs { query, from, to, limit, + cursor, + pages, sort, storage, index, } = args; + if pages == 0 { + bail!("--pages must be at least 1"); + } let api = crate::make_api!(LogsAPI, cfg); let from_ms = util_ext::parse_time_to_unix_millis(&from)?; @@ -232,39 +267,81 @@ pub async fn search(cfg: &Config, args: SearchArgs) -> Result<()> { filter = filter.storage_tier(tier); } - let body = LogsListRequest::new() - .filter(filter) - .page(LogsListRequestPage::new().limit(limit)) - .sort(parse_logs_sort(&sort)); + let mut request_cursor = cursor; + let mut response: Option = None; + let mut next_cursor = None; + + for _ in 0..pages { + let requested_cursor = request_cursor.clone(); + let mut page = LogsListRequestPage::new().limit(limit); + if let Some(cursor) = requested_cursor.clone() { + page = page.cursor(cursor); + } - let params = ListLogsOptionalParams::default().body(body); + let body = LogsListRequest::new() + .filter(filter.clone()) + .page(page) + .sort(parse_logs_sort(&sort)); + let params = ListLogsOptionalParams::default().body(body); + let resp = api + .list_logs(params) + .await + .map_err(|e| anyhow::anyhow!("failed to search logs: {:?}", e))?; + let page_response = serde_json::to_value(&resp)?; + let has_results = page_response + .get("data") + .and_then(|data| data.as_array()) + .is_some_and(|data| !data.is_empty()); + let returned_cursor = next_log_cursor(&page_response); + append_log_page(&mut response, page_response); + next_cursor = returned_cursor.clone(); + + if !has_results { + next_cursor = None; + break; + } - let resp = api - .list_logs(params) - .await - .map_err(|e| anyhow::anyhow!("failed to search logs: {:?}", e))?; + match returned_cursor { + Some(after) if requested_cursor.as_deref() != Some(after.as_str()) => { + request_cursor = Some(after); + } + _ => { + next_cursor = None; + break; + } + } + } + + let response = response.unwrap_or_else(|| serde_json::json!({})); let meta = if cfg.agent_mode { - let count = resp.data.as_ref().map(|d| d.len()); - let truncated = count.is_some_and(|c| c as i32 >= limit); + let count = response + .get("data") + .and_then(|data| data.as_array()) + .map(|data| data.len()); + let truncated = next_cursor.is_some() || count.is_some_and(|c| c as i32 >= limit); Some(formatter::Metadata { count, truncated, command: Some("logs search".into()), - next_action: if truncated { - Some(format!( - "Results may be truncated at {limit}. Use --limit={} or narrow the --query", - limit + 1 - )) - } else { - None - }, + next_action: next_cursor + .map(|cursor| format!("More results available. Use --cursor=\"{cursor}\" to retrieve the next page.")) + .or_else(|| { + if truncated { + Some(format!( + "Results may be truncated at {limit}. Use --limit={} or narrow the --query", + limit + 1 + )) + } else { + None + } + }), }) } else { None }; formatter::format_and_print( - &resp, + &response, &cfg.output_format, cfg.agent_mode, meta.as_ref(), @@ -459,12 +536,41 @@ mod tests { from: "1h".into(), to: "now".into(), limit: 10, + cursor: None, + pages: 1, sort: "-timestamp".into(), storage, index, } } + #[test] + fn test_append_log_page_combines_data_and_keeps_latest_meta() { + let mut response = None; + append_log_page( + &mut response, + serde_json::json!({ + "data": [{"id": "log-1"}], + "meta": {"page": {"after": "cursor-2"}} + }), + ); + append_log_page( + &mut response, + serde_json::json!({ + "data": [{"id": "log-2"}], + "meta": {"page": {}} + }), + ); + + assert_eq!( + response.unwrap(), + serde_json::json!({ + "data": [{"id": "log-1"}, {"id": "log-2"}], + "meta": {"page": {}} + }) + ); + } + #[test] fn test_normalize_storage_tier_alias() { let tier = normalize_storage_tier(Some("online_archives".into())).unwrap(); @@ -817,6 +923,107 @@ mod tests { cleanup_env(); } + #[tokio::test] + async fn test_logs_search_supports_all_output_formats() { + let _lock = lock_env().await; + let mut server = mockito::Server::new_async().await; + let mut cfg = test_config(&server.url()); + let _mock = mock_any( + &mut server, + "POST", + r#"{"data":[{"id":"log-1","attributes":{"message":"error"}}],"meta":{"page":{}}}"#, + ) + .await; + + for format in [ + OutputFormat::Json, + OutputFormat::Yaml, + OutputFormat::Table, + OutputFormat::Csv, + OutputFormat::Tsv, + ] { + cfg.output_format = format.clone(); + let result = super::search(&cfg, search_args("status:error", None, vec![])).await; + assert!( + result.is_ok(), + "logs search failed for {format}: {:?}", + result.err() + ); + } + cleanup_env(); + } + + #[tokio::test] + async fn test_logs_search_with_cursor() { + let _lock = lock_env().await; + let mut server = mockito::Server::new_async().await; + let cfg = test_config(&server.url()); + let _mock = server + .mock("POST", mockito::Matcher::Any) + .match_query(mockito::Matcher::Any) + .match_body(mockito::Matcher::Regex( + r#""cursor":"cursor-abc""#.to_string(), + )) + .with_status(200) + .with_header("content-type", "application/json") + .with_body(r#"{"data": [], "meta": {"page": {}}}"#) + .create_async() + .await; + + let mut args = search_args("status:error", None, vec![]); + args.cursor = Some("cursor-abc".into()); + let result = super::search(&cfg, args).await; + assert!( + result.is_ok(), + "logs search with cursor failed: {:?}", + result.err() + ); + cleanup_env(); + } + + #[tokio::test] + async fn test_logs_search_fetches_requested_pages() { + let _lock = lock_env().await; + let mut server = mockito::Server::new_async().await; + let cfg = test_config(&server.url()); + let first_page = server + .mock("POST", mockito::Matcher::Any) + .match_query(mockito::Matcher::Any) + .match_body(mockito::Matcher::Regex( + r#""page":\{"limit":10\}"#.to_string(), + )) + .with_status(200) + .with_header("content-type", "application/json") + .with_body(r#"{"data":[{"id":"log-1"}],"meta":{"page":{"after":"cursor-2"}}}"#) + .expect(1) + .create_async() + .await; + let second_page = server + .mock("POST", mockito::Matcher::Any) + .match_query(mockito::Matcher::Any) + .match_body(mockito::Matcher::Regex( + r#""cursor":"cursor-2""#.to_string(), + )) + .with_status(200) + .with_header("content-type", "application/json") + .with_body(r#"{"data":[{"id":"log-2"}],"meta":{"page":{"after":"cursor-3"}}}"#) + .expect(1) + .create_async() + .await; + + let mut args = search_args("status:error", None, vec![]); + args.pages = 2; + let result = super::search(&cfg, args).await; + assert!( + result.is_ok(), + "logs search pagination failed: {:?}", + result.err() + ); + first_page.assert_async().await; + second_page.assert_async().await; + cleanup_env(); + } + #[tokio::test] async fn test_logs_search_with_indexes() { let _lock = lock_env().await; diff --git a/src/main.rs b/src/main.rs index cefe69b1..5939c47e 100644 --- a/src/main.rs +++ b/src/main.rs @@ -3146,6 +3146,16 @@ enum LogActions { to: String, #[arg(long, default_value_t = 50, help = "Maximum number of logs (1-1000)")] limit: i32, + #[arg(long, help = "Pagination cursor from a previous response")] + cursor: Option, + #[arg( + long = "max-pages", + visible_alias = "pages", + default_value_t = 1, + value_parser = clap::value_parser!(u32).range(1..), + help = "Maximum pages to fetch automatically (default: 1)" + )] + pages: u32, #[arg(long, help = "Sort order: asc or desc", default_value = "desc")] sort: String, #[arg( @@ -12411,6 +12421,8 @@ async fn main_inner() -> anyhow::Result<()> { from, to, limit, + cursor, + pages, sort, index, storage, @@ -12422,6 +12434,8 @@ async fn main_inner() -> anyhow::Result<()> { from, to, limit, + cursor, + pages, sort, storage, index, @@ -12445,6 +12459,8 @@ async fn main_inner() -> anyhow::Result<()> { from, to, limit, + cursor: None, + pages: 1, sort, storage, index, @@ -12469,6 +12485,8 @@ async fn main_inner() -> anyhow::Result<()> { from, to, limit, + cursor: None, + pages: 1, sort, storage, index, diff --git a/src/test_commands.rs b/src/test_commands.rs index 5fafe698..105bd519 100644 --- a/src/test_commands.rs +++ b/src/test_commands.rs @@ -580,8 +580,11 @@ fn test_logs_search_and_aggregate_default_to_flex_storage() { match search.command { crate::Commands::Logs { - action: crate::LogActions::Search { storage, .. }, - } => assert_eq!(storage.as_deref(), Some("flex")), + action: crate::LogActions::Search { storage, pages, .. }, + } => { + assert_eq!(storage.as_deref(), Some("flex")); + assert_eq!(pages, 1); + } _ => panic!("expected LogActions::Search"), } match aggregate.command { @@ -592,6 +595,44 @@ fn test_logs_search_and_aggregate_default_to_flex_storage() { } } +#[test] +fn test_logs_search_cursor_parses() { + use clap::Parser; + + let cli = crate::Cli::try_parse_from([ + "pup", + "logs", + "search", + "--query", + "*", + "--cursor", + "cursor-abc", + "--pages", + "3", + ]) + .expect("logs search --cursor should parse"); + + match cli.command { + crate::Commands::Logs { + action: crate::LogActions::Search { cursor, pages, .. }, + } => { + assert_eq!(cursor.as_deref(), Some("cursor-abc")); + assert_eq!(pages, 3); + } + _ => panic!("expected LogActions::Search"), + } +} + +#[test] +fn test_logs_search_rejects_zero_pages() { + use clap::Parser; + + let result = + crate::Cli::try_parse_from(["pup", "logs", "search", "--query", "*", "--max-pages", "0"]); + + assert!(result.is_err(), "logs search should reject --max-pages 0"); +} + #[test] fn test_logs_search_and_aggregate_storage_overrides_are_preserved() { use clap::Parser;