1#![forbid(unsafe_code)]
8
9#[doc(hidden)]
10pub mod auth;
11#[doc(hidden)]
12pub mod auth_namespace;
13mod error;
14#[doc(hidden)]
15pub mod flags;
16pub mod help;
17#[doc(hidden)]
18pub mod ids;
19mod parser;
20#[doc(hidden)]
21pub mod render;
22#[doc(hidden)]
23pub mod setup;
24#[doc(hidden)]
25pub mod status;
26#[doc(hidden)]
27pub mod version;
28
29use auth::{open_browser, resolve_credential, write_logout_result};
30pub use error::{CliError, RuntimeError, write_cli_error, write_runtime_error};
31use futures_util::StreamExt as _;
32pub use parser::parse_command;
33use render::write_api_success;
34use setup::write_setup_plan;
35use status::execute_status;
36use tokio_tungstenite::connect_async;
37use tokio_tungstenite::tungstenite::{Error as WebSocketError, Message};
38use version::execute_version;
39
40const WATCH_RECONNECT_INITIAL_DELAY: std::time::Duration = std::time::Duration::from_secs(1);
42const WATCH_RECONNECT_MAX_DELAY: std::time::Duration = std::time::Duration::from_secs(30);
44const WATCH_RECONNECT_JITTER_MAX_MILLIS: u64 = 250;
46
47pub(crate) const ISSUE_STATUS_VALUES_NEXT_STEP: &str =
49 "use one of unresolved/open, resolved/closed, ignored";
50pub(crate) const ISSUE_STATUS_FILTER_NEXT_STEP: &str =
52 "use --status unresolved/open, --status resolved/closed, or --status ignored";
53pub(crate) const ISSUE_STATUS_ARGUMENT_NEXT_STEP: &str =
55 "provide one of unresolved/open, resolved/closed, ignored";
56
57#[derive(Debug, Clone, PartialEq, Eq)]
59pub enum Command {
60 Help {
62 topic: HelpTopic,
64 json: bool,
66 },
67 Login {
69 open_browser: bool,
71 json: bool,
73 },
74 Logout {
76 json: bool,
78 },
79 Setup {
81 auto: bool,
83 yes: bool,
85 json: bool,
87 },
88 Status {
90 json: bool,
92 },
93 Version {
95 json: bool,
97 },
98 Read {
100 target: ReadTarget,
102 options: Box<ReadOptions>,
104 json: bool,
106 },
107 Watch {
109 target: WatchTarget,
111 options: WatchOptions,
113 json: bool,
115 },
116 Explain {
118 target: ExplainTarget,
120 json: bool,
122 },
123 Set {
125 target: SetTarget,
127 json: bool,
129 },
130 ProjectSetupSeen {
132 project_id: String,
134 options: ProjectSetupSeenOptions,
136 json: bool,
138 },
139}
140
141#[derive(Debug, Clone, Copy, PartialEq, Eq)]
143pub enum HelpTopic {
144 Root,
146 Login,
148 Logout,
150 Setup,
152 Status,
154 Version,
156 Auth,
158 Json,
160 Examples,
162 Projects,
164 Usage,
166 Read,
168 ReadLogs,
170 ReadIssues,
172 ReadActions,
174 ReadReleases,
176 ReadTrace,
178 ReadIssue,
180 Watch,
182 Explain,
184 Set,
186}
187
188impl HelpTopic {
189 #[must_use]
191 pub const fn key(self) -> &'static str {
192 match self {
193 Self::Root => "root",
194 Self::Login => "login",
195 Self::Logout => "logout",
196 Self::Setup => "setup",
197 Self::Status => "status",
198 Self::Version => "version",
199 Self::Auth => "auth",
200 Self::Json => "json",
201 Self::Examples => "examples",
202 Self::Projects => "projects",
203 Self::Usage => "usage",
204 Self::Read => "read",
205 Self::ReadLogs => "read_logs",
206 Self::ReadIssues => "read_issues",
207 Self::ReadActions => "read_actions",
208 Self::ReadReleases => "read_releases",
209 Self::ReadTrace => "read_trace",
210 Self::ReadIssue => "read_issue",
211 Self::Watch => "watch",
212 Self::Explain => "explain",
213 Self::Set => "set",
214 }
215 }
216}
217
218#[derive(Debug, Clone, PartialEq, Eq)]
220pub enum ReadTarget {
221 Logs,
223 Issues,
225 Actions,
227 Releases,
229 Trace(String),
231 Issue(String),
233}
234
235#[derive(Debug, Clone, Default, PartialEq, Eq)]
237pub struct ReadOptions {
238 pub name: Option<String>,
240 pub since: Option<String>,
242 pub user: Option<String>,
244 pub trace: Option<String>,
246 pub level: Option<String>,
248 pub search: Option<String>,
250 pub project: Option<String>,
252 pub release: Option<String>,
254 pub environment: Option<String>,
256 pub status: Option<String>,
258 pub limit: Option<String>,
260}
261
262impl ReadOptions {
263 #[must_use]
265 pub(crate) fn first_trace_detail_unsupported_flag(&self) -> Option<&'static str> {
266 first_present_flag([
267 (self.name.is_some(), "--name"),
268 (self.since.is_some(), "--since"),
269 (self.user.is_some(), "--user"),
270 (self.trace.is_some(), "--trace"),
271 (self.level.is_some(), "--severity"),
272 (self.search.is_some(), "--search"),
273 (self.status.is_some(), "--status"),
274 (self.limit.is_some(), "--limit"),
275 ])
276 }
277
278 #[must_use]
280 pub(crate) fn first_issue_detail_unsupported_flag(&self) -> Option<&'static str> {
281 first_present_flag([
282 (self.name.is_some(), "--name"),
283 (self.since.is_some(), "--since"),
284 (self.user.is_some(), "--user"),
285 (self.trace.is_some(), "--trace"),
286 (self.level.is_some(), "--severity"),
287 (self.search.is_some(), "--search"),
288 (self.project.is_some(), "--project"),
289 (self.release.is_some(), "--release"),
290 (self.environment.is_some(), "--environment"),
291 (self.status.is_some(), "--status"),
292 (self.limit.is_some(), "--limit"),
293 ])
294 }
295
296 #[must_use]
298 pub(crate) fn first_log_unsupported_flag(&self) -> Option<&'static str> {
299 first_present_flag([
300 (self.name.is_some(), "--name"),
301 (self.user.is_some(), "--user"),
302 (self.status.is_some(), "--status"),
303 ])
304 }
305
306 #[must_use]
308 pub(crate) fn first_issue_list_unsupported_flag(&self) -> Option<&'static str> {
309 first_present_flag([
310 (self.name.is_some(), "--name"),
311 (self.since.is_some(), "--since"),
312 (self.user.is_some(), "--user"),
313 (self.trace.is_some(), "--trace"),
314 (self.level.is_some(), "--severity"),
315 (self.search.is_some(), "--search"),
316 ])
317 }
318
319 #[must_use]
321 pub(crate) fn first_action_unsupported_flag(&self) -> Option<&'static str> {
322 first_present_flag([
323 (self.trace.is_some(), "--trace"),
324 (self.level.is_some(), "--severity"),
325 (self.search.is_some(), "--search"),
326 (self.status.is_some(), "--status"),
327 ])
328 }
329
330 #[must_use]
332 pub(crate) fn first_release_unsupported_flag(&self) -> Option<&'static str> {
333 first_present_flag([
334 (self.name.is_some(), "--name"),
335 (self.since.is_some(), "--since"),
336 (self.user.is_some(), "--user"),
337 (self.trace.is_some(), "--trace"),
338 (self.level.is_some(), "--severity"),
339 (self.search.is_some(), "--search"),
340 (self.status.is_some(), "--status"),
341 ])
342 }
343}
344
345fn first_present_flag<const N: usize>(flags: [(bool, &'static str); N]) -> Option<&'static str> {
347 flags
348 .iter()
349 .find_map(|(present, flag)| present.then_some(*flag))
350}
351
352#[derive(Debug, Clone, Copy, PartialEq, Eq)]
354pub enum WatchTarget {
355 All,
357 Logs,
359 Issues,
361 Actions,
363}
364
365#[derive(Debug, Clone, Default, PartialEq, Eq)]
367pub struct WatchOptions {
368 pub severity: Vec<String>,
370}
371
372#[derive(Debug, Clone, Default, PartialEq, Eq)]
374pub struct ProjectSetupSeenOptions {
375 pub runtime: Option<String>,
377 pub source: Option<String>,
379 pub environment: Option<String>,
381}
382
383#[derive(Debug, Clone, PartialEq, Eq)]
385pub enum ExplainTarget {
386 Issue(String),
388 Trace(String),
390}
391
392#[derive(Debug, Clone, PartialEq, Eq)]
394pub enum SetTarget {
395 IssueStatus {
397 id: String,
399 status: String,
401 },
402}
403
404#[derive(Debug, Clone, PartialEq, Eq)]
406pub struct CliEnvironment {
407 pub base_url: String,
409 pub token: Option<String>,
411 pub home: Option<std::path::PathBuf>,
413 pub cwd: Option<std::path::PathBuf>,
415}
416
417impl CliEnvironment {
418 #[must_use]
420 pub fn from_process() -> Self {
421 Self {
422 base_url: std::env::var("LOGBREW_API_URL")
423 .unwrap_or_else(|_| String::from("https://api.logbrew.co")),
424 token: std::env::var("LOGBREW_TOKEN").ok(),
425 home: std::env::var_os("HOME").map(std::path::PathBuf::from),
426 cwd: std::env::current_dir().ok(),
427 }
428 }
429}
430
431impl Command {
432 #[must_use]
434 pub fn http_path(&self) -> Option<String> {
435 match self {
436 Self::Read {
437 target, options, ..
438 } => Some(read_path(
439 target,
440 &ReadPathFilters {
441 name: options.name.as_deref(),
442 since: options.since.as_deref(),
443 user: options.user.as_deref(),
444 trace: options.trace.as_deref(),
445 level: options.level.as_deref(),
446 search: options.search.as_deref(),
447 project: options.project.as_deref(),
448 release: options.release.as_deref(),
449 environment: options.environment.as_deref(),
450 status: options.status.as_deref(),
451 limit: options.limit.as_deref(),
452 },
453 )),
454 Self::Explain { target, .. } => Some(explain_path(target)),
455 Self::Set { target, .. } => Some(set_path(target)),
456 Self::ProjectSetupSeen { project_id, .. } => {
457 Some(format!("/api/projects/{project_id}/setup/seen"))
458 }
459 Self::Help { .. }
460 | Self::Login { .. }
461 | Self::Logout { .. }
462 | Self::Setup { .. }
463 | Self::Status { .. }
464 | Self::Version { .. }
465 | Self::Watch { .. } => None,
466 }
467 }
468
469 #[must_use]
471 pub const fn wants_json(&self) -> bool {
472 match self {
473 Self::Help { json, .. }
474 | Self::Login { json, .. }
475 | Self::Logout { json }
476 | Self::Status { json }
477 | Self::Version { json }
478 | Self::Read { json, .. }
479 | Self::Watch { json, .. }
480 | Self::Explain { json, .. }
481 | Self::Set { json, .. }
482 | Self::ProjectSetupSeen { json, .. }
483 | Self::Setup { json, .. } => *json,
484 }
485 }
486
487 #[must_use]
489 pub const fn http_method(&self) -> Option<HttpMethod> {
490 match self {
491 Self::Read { .. } | Self::Explain { .. } => Some(HttpMethod::Get),
492 Self::ProjectSetupSeen { .. } => Some(HttpMethod::Post),
493 Self::Set { .. } => Some(HttpMethod::Patch),
494 Self::Help { .. }
495 | Self::Login { .. }
496 | Self::Logout { .. }
497 | Self::Setup { .. }
498 | Self::Status { .. }
499 | Self::Version { .. }
500 | Self::Watch { .. } => None,
501 }
502 }
503
504 #[must_use]
506 pub fn request_body(&self) -> Option<serde_json::Value> {
507 self.request_body_for_token(None)
508 }
509
510 #[must_use]
512 fn request_body_for_token(&self, token: Option<&str>) -> Option<serde_json::Value> {
513 match self {
514 Self::Set {
515 target: SetTarget::IssueStatus { status, .. },
516 ..
517 } => Some(serde_json::json!({ "status": status })),
518 Self::ProjectSetupSeen { options, .. } => Some(project_setup_seen_body(options, token)),
519 Self::Help { .. }
520 | Self::Login { .. }
521 | Self::Logout { .. }
522 | Self::Setup { .. }
523 | Self::Status { .. }
524 | Self::Version { .. }
525 | Self::Read { .. }
526 | Self::Watch { .. }
527 | Self::Explain { .. } => None,
528 }
529 }
530}
531
532#[derive(Debug, Clone, Copy, PartialEq, Eq)]
534pub enum HttpMethod {
535 Get,
537 Post,
539 Patch,
541}
542
543fn project_setup_seen_body(
545 options: &ProjectSetupSeenOptions,
546 token: Option<&str>,
547) -> serde_json::Value {
548 let mut body = serde_json::Map::new();
549 if let Some(runtime) = options.runtime.as_ref() {
550 drop(body.insert(
551 "runtime".to_owned(),
552 serde_json::Value::String(runtime.clone()),
553 ));
554 }
555 if let Some(source) = setup_seen_source(options, token) {
556 drop(body.insert("source".to_owned(), serde_json::Value::String(source)));
557 }
558 if let Some(environment) = options.environment.as_ref() {
559 drop(body.insert(
560 "environment".to_owned(),
561 serde_json::Value::String(environment.clone()),
562 ));
563 }
564 serde_json::Value::Object(body)
565}
566
567fn setup_seen_source(options: &ProjectSetupSeenOptions, token: Option<&str>) -> Option<String> {
569 if token_is_project_ingest_key(token) {
570 return None;
571 }
572 Some(options.source.as_deref().unwrap_or("cli").to_owned())
573}
574
575fn token_is_project_ingest_key(token: Option<&str>) -> bool {
577 token.is_some_and(|token| token.trim_start().starts_with("lbw_ingest_"))
578}
579
580pub async fn execute_command<W: std::io::Write>(
586 command: &Command,
587 env: &CliEnvironment,
588 output: &mut W,
589) -> Result<(), RuntimeError> {
590 match command {
591 Command::Help { topic, json } => execute_help(*topic, *json, output),
592 Command::Login { open_browser, json } => execute_login(env, *open_browser, *json, output),
593 Command::Logout { json } => execute_logout(env, *json, output),
594 Command::Setup { auto, yes, json } => execute_setup(env, *auto, *yes, *json, output),
595 Command::Status { json } => execute_status(env, *json, output).await,
596 Command::Version { json } => execute_version(*json, output),
597 Command::Read { .. }
598 | Command::Explain { .. }
599 | Command::Set { .. }
600 | Command::ProjectSetupSeen { .. } => execute_http(command, env, output).await,
601 Command::Watch {
602 target,
603 options,
604 json,
605 } => execute_watch(env, *target, options, *json, output).await,
606 }
607}
608
609fn execute_help<W: std::io::Write>(
611 topic: HelpTopic,
612 json: bool,
613 output: &mut W,
614) -> Result<(), RuntimeError> {
615 let help = help::help_text(topic);
616 if json {
617 let body = serde_json::json!({
618 "ok": true,
619 "topic": topic.key(),
620 "help": help,
621 });
622 writeln!(output, "{body}")?;
623 } else {
624 writeln!(output, "{help}")?;
625 }
626 Ok(())
627}
628
629fn execute_login<W: std::io::Write>(
631 env: &CliEnvironment,
632 should_open_browser: bool,
633 json: bool,
634 output: &mut W,
635) -> Result<(), RuntimeError> {
636 let auth_url = format!("{}/api/auth/cli/login", env.base_url.trim_end_matches('/'));
637 let opened = should_open_browser && open_browser(auth_url.as_str());
638
639 if json {
640 let body = serde_json::json!({
641 "ok": true,
642 "auth_url": auth_url,
643 "browser_opened": opened,
644 "next": "open auth_url in a browser",
645 });
646 writeln!(output, "{body}")?;
647 } else {
648 writeln!(output, "Open this URL to log in: {auth_url}")?;
649 writeln!(
650 output,
651 "Browser: {}",
652 if opened { "opened" } else { "not opened" }
653 )?;
654 writeln!(output, "Next: open the URL in a browser")?;
655 }
656 Ok(())
657}
658
659fn execute_logout<W: std::io::Write>(
661 env: &CliEnvironment,
662 json: bool,
663 output: &mut W,
664) -> Result<(), RuntimeError> {
665 write_logout_result(env, json, output)?;
666 Ok(())
667}
668
669fn execute_setup<W: std::io::Write>(
671 env: &CliEnvironment,
672 auto: bool,
673 yes: bool,
674 json: bool,
675 output: &mut W,
676) -> Result<(), RuntimeError> {
677 write_setup_plan(env.cwd.as_deref(), auto, yes, json, output)?;
678 Ok(())
679}
680
681async fn execute_http<W: std::io::Write>(
683 command: &Command,
684 env: &CliEnvironment,
685 output: &mut W,
686) -> Result<(), RuntimeError> {
687 let path = command.http_path().ok_or(CliError::UnknownCommand)?;
688 let url = format!("{}{}", env.base_url.trim_end_matches('/'), path);
689 let client = reqwest::Client::builder()
690 .timeout(std::time::Duration::from_secs(30))
691 .connect_timeout(std::time::Duration::from_secs(10))
692 .build()?;
693
694 let mut request = match command.http_method().unwrap_or(HttpMethod::Get) {
695 HttpMethod::Get => client.get(url),
696 HttpMethod::Post => client.post(url),
697 HttpMethod::Patch => client.patch(url),
698 };
699
700 let credential = resolve_credential(env)?;
701 request = request.bearer_auth(credential.token.as_str());
702
703 if let Some(body) = command.request_body_for_token(Some(credential.token.as_str())) {
704 request = request.json(&body);
705 }
706
707 let response = request.send().await?;
708 let status = response.status();
709 let body = response.text().await?;
710
711 if !status.is_success() {
712 return Err(RuntimeError::Api {
713 status: status.as_u16(),
714 body,
715 auth_source: credential.source,
716 auth_label: credential.label,
717 });
718 }
719
720 write_api_success(command, body.as_str(), output)?;
721 Ok(())
722}
723
724async fn execute_watch<W: std::io::Write>(
726 env: &CliEnvironment,
727 target: WatchTarget,
728 options: &WatchOptions,
729 json: bool,
730 output: &mut W,
731) -> Result<(), RuntimeError> {
732 if !json {
733 return Err(RuntimeError::Unavailable {
734 message: "watch streams JSON for agents",
735 next: "run logbrew watch --json",
736 });
737 }
738
739 let credential = resolve_credential(env)?;
740 let mut reconnect_backoff = WatchReconnectBackoff::default();
741 loop {
742 let ticket = match request_feed_ticket(env, &credential).await {
743 Ok(ticket) => ticket,
744 Err(error) if reconnect_backoff.connected_once() && !runtime_error_is_auth(&error) => {
745 tokio::time::sleep(reconnect_backoff.next_delay()).await;
746 continue;
747 }
748 Err(error) => return Err(error),
749 };
750 let live_url = feed_live_url(env.base_url.as_str(), ticket.as_str())?;
751 let (mut websocket, _) = match connect_async(live_url.as_str()).await {
752 Ok(connection) => connection,
753 Err(error)
754 if reconnect_backoff.connected_once() && !websocket_error_is_auth(&error) =>
755 {
756 tokio::time::sleep(reconnect_backoff.next_delay()).await;
757 continue;
758 }
759 Err(error) => return Err(map_websocket_connect_error(error)),
760 };
761 reconnect_backoff.mark_connected();
762
763 let mut emitted_before_disconnect = false;
764 loop {
765 let Some(message) = websocket.next().await else {
766 break;
767 };
768 let message = match message {
769 Ok(message) => message,
770 Err(error) if websocket_error_is_auth(&error) => {
771 return Err(map_websocket_stream_error(error));
772 }
773 Err(_) => break,
774 };
775 match message {
776 Message::Text(text) => {
777 let event = parse_live_event(text.as_str())?;
778 if watch_event_matches(target, options, &event) {
779 writeln!(output, "{event}")?;
780 }
781 emitted_before_disconnect = true;
782 }
783 Message::Binary(_) | Message::Ping(_) | Message::Pong(_) | Message::Frame(_) => {}
784 Message::Close(_) => return Ok(()),
785 }
786 }
787 if emitted_before_disconnect {
788 reconnect_backoff.reset();
789 }
790 tokio::time::sleep(reconnect_backoff.next_delay()).await;
791 }
792}
793
794#[derive(Debug, Default)]
796struct WatchReconnectBackoff {
797 connected_once: bool,
799 attempts: u32,
801}
802
803impl WatchReconnectBackoff {
804 const fn connected_once(&self) -> bool {
806 self.connected_once
807 }
808
809 const fn mark_connected(&mut self) {
811 self.connected_once = true;
812 }
813
814 const fn reset(&mut self) {
816 self.attempts = 0;
817 }
818
819 fn next_delay(&mut self) -> std::time::Duration {
821 let exponent = self.attempts.min(5);
822 let multiplier = 1_u64 << exponent;
823 self.attempts = self.attempts.saturating_add(1);
824 let base = WATCH_RECONNECT_INITIAL_DELAY
825 .as_secs()
826 .saturating_mul(multiplier)
827 .min(WATCH_RECONNECT_MAX_DELAY.as_secs());
828 let delay = std::time::Duration::from_secs(base) + watch_reconnect_jitter();
829 if delay > WATCH_RECONNECT_MAX_DELAY {
830 WATCH_RECONNECT_MAX_DELAY
831 } else {
832 delay
833 }
834 }
835}
836
837fn watch_reconnect_jitter() -> std::time::Duration {
839 let Ok(elapsed) = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH) else {
840 return std::time::Duration::ZERO;
841 };
842 std::time::Duration::from_millis(
843 u64::from(elapsed.subsec_millis()) % WATCH_RECONNECT_JITTER_MAX_MILLIS,
844 )
845}
846
847const fn runtime_error_is_auth(error: &RuntimeError) -> bool {
849 matches!(
850 error,
851 RuntimeError::MissingToken | RuntimeError::Api { status: 401, .. }
852 )
853}
854
855fn websocket_error_is_auth(error: &WebSocketError) -> bool {
857 matches!(error, WebSocketError::Http(response) if response.status().as_u16() == 401)
858}
859
860async fn request_feed_ticket(
862 env: &CliEnvironment,
863 credential: &auth::AuthCredential,
864) -> Result<String, RuntimeError> {
865 let url = format!("{}/api/feed/ticket", env.base_url.trim_end_matches('/'));
866 let client = reqwest::Client::builder()
867 .timeout(std::time::Duration::from_secs(30))
868 .connect_timeout(std::time::Duration::from_secs(10))
869 .build()?;
870 let response = client
871 .post(url)
872 .bearer_auth(credential.token.as_str())
873 .send()
874 .await?;
875 let status = response.status();
876 let body = response.text().await?;
877 if !status.is_success() {
878 return Err(RuntimeError::Api {
879 status: status.as_u16(),
880 body,
881 auth_source: credential.source,
882 auth_label: credential.label,
883 });
884 }
885
886 let value = serde_json::from_str::<serde_json::Value>(body.as_str()).map_err(|_| {
887 RuntimeError::Unavailable {
888 message: "feed ticket response was not valid JSON",
889 next: "retry logbrew watch or run logbrew status",
890 }
891 })?;
892 value
893 .get("ticket")
894 .and_then(serde_json::Value::as_str)
895 .map(str::trim)
896 .filter(|ticket| !ticket.is_empty())
897 .map(ToOwned::to_owned)
898 .ok_or(RuntimeError::Unavailable {
899 message: "feed ticket response did not include a ticket",
900 next: "retry logbrew watch or run logbrew status",
901 })
902}
903
904fn feed_live_url(base_url: &str, ticket: &str) -> Result<String, RuntimeError> {
906 let trimmed = base_url.trim_end_matches('/');
907 let (scheme, rest) = websocket_base_parts(trimmed).ok_or(RuntimeError::Unavailable {
908 message: "LOGBREW_API_URL must start with http:// or https://",
909 next: "check LOGBREW_API_URL or run logbrew status",
910 })?;
911 Ok(format!(
912 "{scheme}://{rest}/api/feed/live?ticket={}",
913 encode_component(ticket)
914 ))
915}
916
917fn websocket_base_parts(base_url: &str) -> Option<(&'static str, &str)> {
919 base_url
920 .strip_prefix("https://")
921 .map(|rest| ("wss", rest))
922 .or_else(|| base_url.strip_prefix("http://").map(|rest| ("ws", rest)))
923}
924
925fn parse_live_event(text: &str) -> Result<serde_json::Value, RuntimeError> {
927 serde_json::from_str::<serde_json::Value>(text).map_err(|_| RuntimeError::Unavailable {
928 message: "live watch event was not valid JSON",
929 next: "retry logbrew watch or check LOGBREW_API_URL",
930 })
931}
932
933fn watch_event_matches(
935 target: WatchTarget,
936 options: &WatchOptions,
937 event: &serde_json::Value,
938) -> bool {
939 target_matches_event(target, event) && severity_matches(options, event)
940}
941
942fn target_matches_event(target: WatchTarget, event: &serde_json::Value) -> bool {
944 let event_type = event
945 .get("type")
946 .and_then(serde_json::Value::as_str)
947 .unwrap_or_default();
948 match target {
949 WatchTarget::All => true,
950 WatchTarget::Logs => event_type == "native_log",
951 WatchTarget::Issues => event_type == "native_issue",
952 WatchTarget::Actions => event_type == "native_action",
953 }
954}
955
956fn severity_matches(options: &WatchOptions, event: &serde_json::Value) -> bool {
958 if options.severity.is_empty() {
959 return true;
960 }
961 let Some(severity) = event
962 .get("data")
963 .and_then(|data| data.get("severity").or_else(|| data.get("level")))
964 .and_then(serde_json::Value::as_str)
965 else {
966 return false;
967 };
968 options
969 .severity
970 .iter()
971 .any(|allowed| allowed.as_str() == severity)
972}
973
974fn map_websocket_connect_error(error: WebSocketError) -> RuntimeError {
976 match error {
977 WebSocketError::Http(response) if response.status().as_u16() == 401 => {
978 RuntimeError::Unavailable {
979 message: "live watch ticket was rejected",
980 next: "run logbrew login",
981 }
982 }
983 WebSocketError::Http(_) => RuntimeError::Unavailable {
984 message: "live watch websocket upgrade failed",
985 next: "retry logbrew watch or check LOGBREW_API_URL",
986 },
987 WebSocketError::ConnectionClosed
988 | WebSocketError::AlreadyClosed
989 | WebSocketError::Io(_)
990 | WebSocketError::Tls(_)
991 | WebSocketError::Capacity(_)
992 | WebSocketError::Protocol(_)
993 | WebSocketError::WriteBufferFull(_)
994 | WebSocketError::Utf8(_)
995 | WebSocketError::AttackAttempt
996 | WebSocketError::Url(_)
997 | WebSocketError::HttpFormat(_) => RuntimeError::Unavailable {
998 message: "live watch websocket failed",
999 next: "retry logbrew watch or check LOGBREW_API_URL",
1000 },
1001 }
1002}
1003
1004fn map_websocket_stream_error(error: WebSocketError) -> RuntimeError {
1006 match error {
1007 WebSocketError::ConnectionClosed | WebSocketError::AlreadyClosed => {
1008 RuntimeError::Unavailable {
1009 message: "live watch websocket closed",
1010 next: "retry logbrew watch",
1011 }
1012 }
1013 WebSocketError::Http(response) if response.status().as_u16() == 401 => {
1014 RuntimeError::Unavailable {
1015 message: "live watch ticket was rejected",
1016 next: "run logbrew login",
1017 }
1018 }
1019 WebSocketError::Http(_)
1020 | WebSocketError::Io(_)
1021 | WebSocketError::Tls(_)
1022 | WebSocketError::Capacity(_)
1023 | WebSocketError::Protocol(_)
1024 | WebSocketError::WriteBufferFull(_)
1025 | WebSocketError::Utf8(_)
1026 | WebSocketError::AttackAttempt
1027 | WebSocketError::Url(_)
1028 | WebSocketError::HttpFormat(_) => RuntimeError::Unavailable {
1029 message: "live watch websocket failed",
1030 next: "retry logbrew watch or check LOGBREW_API_URL",
1031 },
1032 }
1033}
1034
1035struct ReadPathFilters<'a> {
1037 name: Option<&'a str>,
1039 since: Option<&'a str>,
1041 user: Option<&'a str>,
1043 trace: Option<&'a str>,
1045 level: Option<&'a str>,
1047 search: Option<&'a str>,
1049 project: Option<&'a str>,
1051 release: Option<&'a str>,
1053 environment: Option<&'a str>,
1055 status: Option<&'a str>,
1057 limit: Option<&'a str>,
1059}
1060
1061fn read_path(target: &ReadTarget, filters: &ReadPathFilters<'_>) -> String {
1063 match target {
1064 ReadTarget::Logs => path_with_query(
1065 "/api/logs",
1066 &[
1067 ("severity", filters.level),
1068 ("search", filters.search),
1069 ("since", filters.since),
1070 ("trace_id", filters.trace),
1071 ("project_id", filters.project),
1072 ("release", filters.release),
1073 ("environment", filters.environment),
1074 ("limit", filters.limit),
1075 ],
1076 ),
1077 ReadTarget::Issues => path_with_query(
1078 "/api/telemetry/issues",
1079 &[
1080 ("status", filters.status),
1081 ("project_id", filters.project),
1082 ("release", filters.release),
1083 ("environment", filters.environment),
1084 ("limit", filters.limit),
1085 ],
1086 ),
1087 ReadTarget::Actions => path_with_query(
1088 "/api/telemetry/actions",
1089 &[
1090 ("name", filters.name),
1091 ("since", filters.since),
1092 ("distinct_id", filters.user),
1093 ("project_id", filters.project),
1094 ("release", filters.release),
1095 ("environment", filters.environment),
1096 ("limit", filters.limit),
1097 ],
1098 ),
1099 ReadTarget::Releases => path_with_query(
1100 "/api/telemetry/releases",
1101 &[
1102 ("project_id", filters.project),
1103 ("release", filters.release),
1104 ("environment", filters.environment),
1105 ("limit", filters.limit),
1106 ],
1107 ),
1108 ReadTarget::Trace(id) => path_with_query(
1109 &format!("/api/telemetry/traces/{}", encode_component(id)),
1110 &[
1111 ("project_id", filters.project),
1112 ("release", filters.release),
1113 ("environment", filters.environment),
1114 ],
1115 ),
1116 ReadTarget::Issue(id) => format!("/api/telemetry/issues/{}", encode_component(id)),
1117 }
1118}
1119
1120fn explain_path(target: &ExplainTarget) -> String {
1122 match target {
1123 ExplainTarget::Issue(id) => format!("/api/telemetry/issues/{}", encode_component(id)),
1124 ExplainTarget::Trace(id) => format!("/api/telemetry/traces/{}", encode_component(id)),
1125 }
1126}
1127
1128fn set_path(target: &SetTarget) -> String {
1130 match target {
1131 SetTarget::IssueStatus { id, .. } => {
1132 format!("/api/telemetry/issues/{}", encode_component(id))
1133 }
1134 }
1135}
1136
1137fn path_with_query(path: &str, params: &[(&str, Option<&str>)]) -> String {
1139 let query = params
1140 .iter()
1141 .filter_map(|(name, value)| value.map(|v| format!("{name}={}", encode_component(v))))
1142 .collect::<Vec<_>>();
1143
1144 if query.is_empty() {
1145 path.to_owned()
1146 } else {
1147 format!("{path}?{}", query.join("&"))
1148 }
1149}
1150
1151fn encode_component(value: &str) -> String {
1153 let mut encoded = String::new();
1154 for byte in value.bytes() {
1155 if byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.' | b'~') {
1156 encoded.push(char::from(byte));
1157 } else {
1158 encoded.push('%');
1159 encoded.push(hex_digit(byte >> 4));
1160 encoded.push(hex_digit(byte & 0x0f));
1161 }
1162 }
1163 encoded
1164}
1165
1166fn hex_digit(nibble: u8) -> char {
1168 match nibble {
1169 0..=9 => char::from(b'0' + nibble),
1170 10..=15 => char::from(b'A' + (nibble - 10)),
1171 _ => '?',
1172 }
1173}