#![forbid(unsafe_code)]
#[doc(hidden)]
pub mod auth;
#[doc(hidden)]
pub mod auth_namespace;
#[doc(hidden)]
pub mod doctor;
mod error;
#[doc(hidden)]
pub mod flags;
pub mod help;
#[doc(hidden)]
pub mod ids;
#[doc(hidden)]
pub mod investigate;
mod native_debug_artifacts;
mod parser;
mod project_create;
mod projects;
#[doc(hidden)]
pub mod render;
#[doc(hidden)]
pub mod setup;
#[doc(hidden)]
pub mod status;
mod support;
mod usage;
#[doc(hidden)]
pub mod version;
use auth::{
AuthCredential, execute_login, execute_logout, execute_whoami, send_authenticated_with_refresh,
token_is_project_ingest_key,
};
pub use error::{
CliError, RuntimeError, write_cli_error, write_native_debug_runtime_error, write_runtime_error,
};
use futures_util::StreamExt as _;
pub use parser::parse_command;
use render::write_api_success;
use setup::write_setup_plan;
use status::execute_status;
use tokio_tungstenite::connect_async;
use tokio_tungstenite::tungstenite::{Error as WebSocketError, Message};
use version::execute_version;
const WATCH_RECONNECT_INITIAL_DELAY: std::time::Duration = std::time::Duration::from_secs(1);
const WATCH_RECONNECT_MAX_DELAY: std::time::Duration = std::time::Duration::from_secs(30);
const WATCH_RECONNECT_JITTER_MAX_MILLIS: u64 = 250;
pub(crate) const ISSUE_STATUS_VALUES_NEXT_STEP: &str =
"use one of unresolved/open, resolved/closed, ignored";
pub(crate) const ISSUE_STATUS_FILTER_NEXT_STEP: &str =
"use --status unresolved/open, --status resolved/closed, or --status ignored";
pub(crate) const ISSUE_STATUS_ARGUMENT_NEXT_STEP: &str =
"provide one of unresolved/open, resolved/closed, ignored";
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub enum LoginProvider {
#[default]
GitHub,
GitLab,
Bitbucket,
}
impl LoginProvider {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::GitHub => "github",
Self::GitLab => "gitlab",
Self::Bitbucket => "bitbucket",
}
}
}
impl std::str::FromStr for LoginProvider {
type Err = CliError;
fn from_str(value: &str) -> Result<Self, Self::Err> {
match value {
"github" => Ok(Self::GitHub),
"gitlab" => Ok(Self::GitLab),
"bitbucket" => Ok(Self::Bitbucket),
_ => Err(CliError::InvalidLoginProvider),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Command {
Help {
topic: HelpTopic,
json: bool,
},
Login {
provider: LoginProvider,
open_browser: bool,
json: bool,
},
Logout {
json: bool,
},
Setup {
auto: bool,
yes: bool,
json: bool,
},
Status {
json: bool,
},
WhoAmI {
json: bool,
},
Doctor {
project_id: String,
json: bool,
},
ProjectCreate {
options: ProjectCreateOptions,
json: bool,
},
Projects {
json: bool,
},
Usage {
json: bool,
},
Version {
json: bool,
},
Read {
target: ReadTarget,
options: Box<ReadOptions>,
json: bool,
},
Watch {
target: WatchTarget,
options: WatchOptions,
json: bool,
},
Explain {
target: ExplainTarget,
json: bool,
},
InvestigateIssue {
issue_id: String,
json: bool,
},
NativeDebugArtifacts {
target: NativeDebugArtifactsTarget,
json: bool,
},
Set {
target: SetTarget,
json: bool,
},
ProjectSetupSeen {
project_id: String,
options: ProjectSetupSeenOptions,
json: bool,
},
Support {
target: SupportTarget,
json: bool,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum HelpTopic {
Root,
Login,
Logout,
Setup,
Status,
Version,
Auth,
Json,
Examples,
Projects,
Usage,
Read,
ReadLogs,
ReadIssues,
ReadActions,
ReadReleases,
ReadTraces,
ReadTrace,
ReadIssue,
Watch,
Explain,
Investigate,
NativeDebugArtifacts,
Set,
Support,
}
impl HelpTopic {
#[must_use]
pub const fn key(self) -> &'static str {
match self {
Self::Root => "root",
Self::Login => "login",
Self::Logout => "logout",
Self::Setup => "setup",
Self::Status => "status",
Self::Version => "version",
Self::Auth => "auth",
Self::Json => "json",
Self::Examples => "examples",
Self::Projects => "projects",
Self::Usage => "usage",
Self::Read => "read",
Self::ReadLogs => "read_logs",
Self::ReadIssues => "read_issues",
Self::ReadActions => "read_actions",
Self::ReadReleases => "read_releases",
Self::ReadTraces => "read_traces",
Self::ReadTrace => "read_trace",
Self::ReadIssue => "read_issue",
Self::Watch => "watch",
Self::Explain => "explain",
Self::Investigate => "investigate",
Self::NativeDebugArtifacts => "debug_artifacts",
Self::Set => "set",
Self::Support => "support",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ReadTarget {
Logs,
Issues,
Actions,
Releases,
Traces,
Trace(String),
Issue(String),
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ReadOptions {
pub name: Option<String>,
pub service: Option<String>,
pub since: Option<String>,
pub user: Option<String>,
pub trace: Option<String>,
pub level: Option<String>,
pub search: Option<String>,
pub project: Option<String>,
pub release: Option<String>,
pub environment: Option<String>,
pub status: Option<String>,
pub limit: Option<String>,
pub min_duration_ms: Option<String>,
pub pagination: Option<String>,
pub cursor_time: Option<String>,
pub cursor_id: Option<String>,
}
impl ReadOptions {
#[must_use]
pub(crate) fn first_trace_detail_unsupported_flag(&self) -> Option<&'static str> {
first_present_flag([
(self.name.is_some(), "--name"),
(self.service.is_some(), "--service"),
(self.since.is_some(), "--since"),
(self.user.is_some(), "--user"),
(self.trace.is_some(), "--trace"),
(self.level.is_some(), "--severity"),
(self.search.is_some(), "--search"),
(self.status.is_some(), "--status"),
(self.limit.is_some(), "--limit"),
(self.min_duration_ms.is_some(), "--min-duration-ms"),
(self.pagination.is_some(), "--pagination"),
(self.cursor_time.is_some(), "--cursor-time"),
(self.cursor_id.is_some(), "--cursor-id"),
])
}
#[must_use]
pub(crate) fn first_issue_detail_unsupported_flag(&self) -> Option<&'static str> {
first_present_flag([
(self.name.is_some(), "--name"),
(self.service.is_some(), "--service"),
(self.since.is_some(), "--since"),
(self.user.is_some(), "--user"),
(self.trace.is_some(), "--trace"),
(self.level.is_some(), "--severity"),
(self.search.is_some(), "--search"),
(self.project.is_some(), "--project"),
(self.release.is_some(), "--release"),
(self.environment.is_some(), "--environment"),
(self.status.is_some(), "--status"),
(self.limit.is_some(), "--limit"),
(self.min_duration_ms.is_some(), "--min-duration-ms"),
(self.pagination.is_some(), "--pagination"),
(self.cursor_time.is_some(), "--cursor-time"),
(self.cursor_id.is_some(), "--cursor-id"),
])
}
#[must_use]
pub(crate) fn first_log_unsupported_flag(&self) -> Option<&'static str> {
first_present_flag([
(self.name.is_some(), "--name"),
(self.user.is_some(), "--user"),
(self.status.is_some(), "--status"),
(self.min_duration_ms.is_some(), "--min-duration-ms"),
])
}
#[must_use]
pub(crate) fn first_issue_list_unsupported_flag(&self) -> Option<&'static str> {
first_present_flag([
(self.name.is_some(), "--name"),
(self.user.is_some(), "--user"),
(self.trace.is_some(), "--trace"),
(self.level.is_some(), "--severity"),
(self.search.is_some(), "--search"),
(self.min_duration_ms.is_some(), "--min-duration-ms"),
])
}
#[must_use]
pub(crate) fn first_action_unsupported_flag(&self) -> Option<&'static str> {
first_present_flag([
(self.trace.is_some(), "--trace"),
(self.level.is_some(), "--severity"),
(self.search.is_some(), "--search"),
(self.status.is_some(), "--status"),
(self.min_duration_ms.is_some(), "--min-duration-ms"),
])
}
#[must_use]
pub(crate) fn first_release_unsupported_flag(&self) -> Option<&'static str> {
first_present_flag([
(self.name.is_some(), "--name"),
(self.user.is_some(), "--user"),
(self.trace.is_some(), "--trace"),
(self.level.is_some(), "--severity"),
(self.search.is_some(), "--search"),
(self.status.is_some(), "--status"),
(self.min_duration_ms.is_some(), "--min-duration-ms"),
(self.pagination.is_some(), "--pagination"),
(self.cursor_time.is_some(), "--cursor-time"),
(self.cursor_id.is_some(), "--cursor-id"),
])
}
#[must_use]
pub(crate) fn first_trace_list_unsupported_flag(&self) -> Option<&'static str> {
first_present_flag([
(self.name.is_some(), "--name"),
(self.user.is_some(), "--user"),
(self.trace.is_some(), "--trace"),
(self.level.is_some(), "--severity"),
(self.search.is_some(), "--search"),
(self.pagination.is_some(), "--pagination"),
(self.cursor_time.is_some(), "--cursor-time"),
(self.cursor_id.is_some(), "--cursor-id"),
])
}
}
fn first_present_flag<const N: usize>(flags: [(bool, &'static str); N]) -> Option<&'static str> {
flags
.iter()
.find_map(|(present, flag)| present.then_some(*flag))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum WatchTarget {
All,
Logs,
Issues,
Actions,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct WatchOptions {
pub severity: Vec<String>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ProjectSetupSeenOptions {
pub runtime: Option<String>,
pub source: Option<String>,
pub environment: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProjectCreateOptions {
pub name: String,
pub runtime: Option<String>,
pub environment: Option<String>,
pub ingest_key_file: String,
pub abandon_retry: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum NativeDebugArtifactsTarget {
Upload(NativeDebugUploadOptions),
Lookup(NativeDebugLookupOptions),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NativeDebugLookupOptions {
pub project_id: String,
pub release: String,
pub environment: String,
pub service: String,
pub image_uuid: String,
pub architecture: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NativeDebugUploadOptions {
pub path: String,
pub project_id: String,
pub release: String,
pub environment: String,
pub service: String,
pub expected_image_uuids: Vec<String>,
pub dry_run: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SupportTarget {
Create(Box<SupportTicketCreateOptions>),
List(Box<SupportTicketListOptions>),
Detail(String),
ContextHistory {
ticket_id: String,
},
ReplyContext(Box<SupportContextReplyOptions>),
UpdateStatus {
ticket_id: String,
status: SupportTicketLifecycleStatus,
},
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct SupportContextReplyOptions {
pub ticket_id: String,
pub context: String,
pub retry_key: String,
pub diagnostics: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SupportTicketLifecycleStatus {
Open,
Closed,
}
impl SupportTicketLifecycleStatus {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Open => "open",
Self::Closed => "closed",
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct SupportTicketCreateOptions {
pub category: String,
pub title: String,
pub description: String,
pub project_id: Option<String>,
pub environment: Option<String>,
pub runtime: Option<String>,
pub framework: Option<String>,
pub sdk_package: Option<String>,
pub sdk_version: Option<String>,
pub release: Option<String>,
pub trace_id: Option<String>,
pub event_id: Option<String>,
pub diagnostics: bool,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct SupportTicketListOptions {
pub project_id: Option<String>,
pub status: Option<String>,
pub source: Option<String>,
pub category: Option<String>,
pub release: Option<String>,
pub limit: Option<String>,
pub pagination: Option<String>,
pub cursor_time: Option<String>,
pub cursor_id: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ExplainTarget {
Issue(String),
Trace(String),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SetTarget {
IssueStatus {
id: String,
status: String,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CliEnvironment {
pub base_url: String,
pub token: Option<String>,
pub home: Option<std::path::PathBuf>,
pub cwd: Option<std::path::PathBuf>,
}
impl CliEnvironment {
#[must_use]
pub fn from_process() -> Self {
Self {
base_url: std::env::var("LOGBREW_API_URL")
.unwrap_or_else(|_| String::from("https://api.logbrew.co")),
token: std::env::var("LOGBREW_TOKEN").ok(),
home: std::env::var_os("HOME").map(std::path::PathBuf::from),
cwd: std::env::current_dir().ok(),
}
}
}
impl Command {
#[must_use]
pub fn http_path(&self) -> Option<String> {
match self {
Self::Read {
target, options, ..
} => Some(read_path(
target,
&ReadPathFilters {
name: options.name.as_deref(),
service: options.service.as_deref(),
since: options.since.as_deref(),
user: options.user.as_deref(),
trace: options.trace.as_deref(),
level: options.level.as_deref(),
search: options.search.as_deref(),
project: options.project.as_deref(),
release: options.release.as_deref(),
environment: options.environment.as_deref(),
status: options.status.as_deref(),
limit: options.limit.as_deref(),
min_duration_ms: options.min_duration_ms.as_deref(),
pagination: options.pagination.as_deref(),
cursor_time: options.cursor_time.as_deref(),
cursor_id: options.cursor_id.as_deref(),
},
)),
Self::Explain { target, .. } => Some(explain_path(target)),
Self::Set { target, .. } => Some(set_path(target)),
Self::ProjectSetupSeen { project_id, .. } => {
Some(format!("/api/projects/{project_id}/setup/seen"))
}
Self::ProjectCreate { .. } | Self::Projects { .. } => {
Some(String::from("/api/projects"))
}
Self::Support { target, .. } => Some(support::path(target)),
Self::Help { .. }
| Self::Login { .. }
| Self::Logout { .. }
| Self::Setup { .. }
| Self::Status { .. }
| Self::WhoAmI { .. }
| Self::Doctor { .. }
| Self::Usage { .. }
| Self::Version { .. }
| Self::InvestigateIssue { .. }
| Self::NativeDebugArtifacts { .. }
| Self::Watch { .. } => None,
}
}
#[must_use]
pub const fn wants_json(&self) -> bool {
match self {
Self::Help { json, .. }
| Self::Login { json, .. }
| Self::Logout { json }
| Self::Status { json }
| Self::WhoAmI { json }
| Self::Doctor { json, .. }
| Self::ProjectCreate { json, .. }
| Self::Projects { json }
| Self::Usage { json }
| Self::Version { json }
| Self::Read { json, .. }
| Self::Watch { json, .. }
| Self::Explain { json, .. }
| Self::InvestigateIssue { json, .. }
| Self::NativeDebugArtifacts { json, .. }
| Self::Set { json, .. }
| Self::ProjectSetupSeen { json, .. }
| Self::Support { json, .. }
| Self::Setup { json, .. } => *json,
}
}
#[must_use]
pub const fn http_method(&self) -> Option<HttpMethod> {
match self {
Self::ProjectCreate { .. }
| Self::ProjectSetupSeen { .. }
| Self::Support {
target: SupportTarget::Create(_),
..
}
| Self::Support {
target: SupportTarget::ReplyContext(_),
..
} => Some(HttpMethod::Post),
Self::Support {
target: SupportTarget::UpdateStatus { .. },
..
}
| Self::Set { .. } => Some(HttpMethod::Patch),
Self::Projects { .. }
| Self::Read { .. }
| Self::Explain { .. }
| Self::Support { .. } => Some(HttpMethod::Get),
Self::Help { .. }
| Self::Login { .. }
| Self::Logout { .. }
| Self::Setup { .. }
| Self::Status { .. }
| Self::WhoAmI { .. }
| Self::Doctor { .. }
| Self::Usage { .. }
| Self::Version { .. }
| Self::InvestigateIssue { .. }
| Self::NativeDebugArtifacts { .. }
| Self::Watch { .. } => None,
}
}
#[must_use]
pub fn request_body(&self) -> Option<serde_json::Value> {
self.request_body_for_token(None)
}
#[must_use]
fn request_body_for_token(&self, token: Option<&str>) -> Option<serde_json::Value> {
match self {
Self::Set {
target: SetTarget::IssueStatus { status, .. },
..
} => Some(serde_json::json!({ "status": status })),
Self::ProjectSetupSeen { options, .. } => Some(project_setup_seen_body(options, token)),
Self::ProjectCreate { options, .. } => Some(project_create_body(options)),
Self::Support {
target: SupportTarget::Create(options),
..
} => Some(support::create_body(options)),
Self::Support {
target: SupportTarget::ReplyContext(options),
..
} => Some(support::context_body(options)),
Self::Support {
target: SupportTarget::UpdateStatus { status, .. },
..
} => Some(serde_json::json!({"status": status.as_str()})),
Self::Help { .. }
| Self::Login { .. }
| Self::Logout { .. }
| Self::Setup { .. }
| Self::Status { .. }
| Self::WhoAmI { .. }
| Self::Doctor { .. }
| Self::Projects { .. }
| Self::Usage { .. }
| Self::Version { .. }
| Self::Read { .. }
| Self::Watch { .. }
| Self::Explain { .. }
| Self::InvestigateIssue { .. }
| Self::NativeDebugArtifacts { .. }
| Self::Support { .. } => None,
}
}
fn idempotency_key(&self) -> Option<&str> {
match self {
Self::Support {
target: SupportTarget::ReplyContext(options),
..
} => Some(options.retry_key.as_str()),
Self::Help { .. }
| Self::Login { .. }
| Self::Logout { .. }
| Self::Setup { .. }
| Self::Status { .. }
| Self::WhoAmI { .. }
| Self::Doctor { .. }
| Self::Usage { .. }
| Self::Version { .. }
| Self::Read { .. }
| Self::Watch { .. }
| Self::Explain { .. }
| Self::InvestigateIssue { .. }
| Self::NativeDebugArtifacts { .. }
| Self::Set { .. }
| Self::ProjectSetupSeen { .. }
| Self::ProjectCreate { .. }
| Self::Projects { .. }
| Self::Support { .. } => None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum HttpMethod {
Get,
Post,
Patch,
}
fn project_setup_seen_body(
options: &ProjectSetupSeenOptions,
token: Option<&str>,
) -> serde_json::Value {
let mut body = serde_json::Map::new();
if let Some(runtime) = options.runtime.as_ref() {
drop(body.insert(
"runtime".to_owned(),
serde_json::Value::String(runtime.clone()),
));
}
if let Some(source) = setup_seen_source(options, token) {
drop(body.insert("source".to_owned(), serde_json::Value::String(source)));
}
if let Some(environment) = options.environment.as_ref() {
drop(body.insert(
"environment".to_owned(),
serde_json::Value::String(environment.clone()),
));
}
serde_json::Value::Object(body)
}
fn project_create_body(options: &ProjectCreateOptions) -> serde_json::Value {
let mut body = serde_json::Map::new();
drop(body.insert(
"name".to_owned(),
serde_json::Value::String(options.name.clone()),
));
if let Some(runtime) = options.runtime.as_ref() {
drop(body.insert(
"runtime".to_owned(),
serde_json::Value::String(runtime.clone()),
));
}
if let Some(environment) = options.environment.as_ref() {
drop(body.insert(
"environment".to_owned(),
serde_json::Value::String(environment.clone()),
));
}
drop(body.insert(
"source".to_owned(),
serde_json::Value::String(String::from("cli")),
));
serde_json::Value::Object(body)
}
fn setup_seen_source(options: &ProjectSetupSeenOptions, token: Option<&str>) -> Option<String> {
if token_is_project_ingest_key(token) {
return None;
}
Some(options.source.as_deref().unwrap_or("cli").to_owned())
}
pub async fn execute_command<W: std::io::Write>(
command: &Command,
env: &CliEnvironment,
output: &mut W,
) -> Result<(), RuntimeError> {
match command {
Command::Help { topic, json } => execute_help(*topic, *json, output),
Command::Login {
provider,
open_browser,
json,
} => execute_login(env, *provider, *open_browser, *json, output).await,
Command::Logout { json } => execute_logout(env, *json, output).await,
Command::Setup { auto, yes, json } => execute_setup(env, *auto, *yes, *json, output),
Command::Status { json } => execute_status(env, *json, output).await,
Command::WhoAmI { json } => execute_whoami(env, *json, output).await,
Command::Doctor { project_id, json } => {
doctor::execute(env, project_id.as_str(), *json, output).await
}
Command::ProjectCreate { options, json } => {
project_create::execute(env, options, *json, output).await
}
Command::Projects { json } => projects::execute(env, *json, output).await,
Command::Usage { json } => usage::execute(env, *json, output).await,
Command::Version { json } => execute_version(*json, output),
Command::InvestigateIssue { issue_id, json } => {
investigate::execute(env, issue_id.as_str(), *json, output).await
}
Command::NativeDebugArtifacts { target, json } => {
native_debug_artifacts::execute(env, target, *json, output).await
}
Command::Read { .. }
| Command::Explain { .. }
| Command::Set { .. }
| Command::ProjectSetupSeen { .. }
| Command::Support { .. } => execute_http(command, env, output).await,
Command::Watch {
target,
options,
json,
} => execute_watch(env, *target, options, *json, output).await,
}
}
fn execute_help<W: std::io::Write>(
topic: HelpTopic,
json: bool,
output: &mut W,
) -> Result<(), RuntimeError> {
let help = help::help_text(topic);
if json {
let body = serde_json::json!({
"ok": true,
"topic": topic.key(),
"help": help,
});
writeln!(output, "{body}")?;
} else {
writeln!(output, "{help}")?;
}
Ok(())
}
fn execute_setup<W: std::io::Write>(
env: &CliEnvironment,
auto: bool,
yes: bool,
json: bool,
output: &mut W,
) -> Result<(), RuntimeError> {
write_setup_plan(env.cwd.as_deref(), auto, yes, json, output)?;
Ok(())
}
async fn execute_http<W: std::io::Write>(
command: &Command,
env: &CliEnvironment,
output: &mut W,
) -> Result<(), RuntimeError> {
let path = command.http_path().ok_or(CliError::UnknownCommand)?;
let url = format!("{}{}", env.base_url.trim_end_matches('/'), path);
let client = reqwest::Client::builder()
.timeout(std::time::Duration::from_secs(30))
.connect_timeout(std::time::Duration::from_secs(10))
.build()?;
let support_command = matches!(command, Command::Support { .. });
let response_result = send_authenticated_with_refresh(&client, env, |client, credential| {
build_command_request(client, command, url.as_str(), credential)
})
.await;
let (response, credential) = match response_result {
Ok(response) => response,
Err(RuntimeError::Http(_)) if support_command => return Err(support_transport_error()),
Err(error) => return Err(error),
};
let status = response.status();
let body = match response.text().await {
Ok(body) => body,
Err(_) if support_command => return Err(support_transport_error()),
Err(error) => return Err(RuntimeError::Http(error)),
};
if !status.is_success() {
let body = if let Command::Support { target, .. } = command {
support::safe_error_body(target, status.as_u16())
} else {
credential.redact_response_body(body.as_str())
};
return Err(RuntimeError::Api {
status: status.as_u16(),
body,
auth_source: credential.source(),
auth_label: credential.label(),
});
}
write_api_success(command, body.as_str(), output)?;
Ok(())
}
const fn support_transport_error() -> RuntimeError {
RuntimeError::Unavailable {
message: "support request could not be completed",
next: "check network connectivity and retry the support command",
}
}
fn build_command_request(
client: &reqwest::Client,
command: &Command,
url: &str,
credential: &AuthCredential,
) -> reqwest::RequestBuilder {
let mut request = match command.http_method().unwrap_or(HttpMethod::Get) {
HttpMethod::Get => client.get(url),
HttpMethod::Post => client.post(url),
HttpMethod::Patch => client.patch(url),
}
.bearer_auth(credential.token());
if let Some(body) = command.request_body_for_token(Some(credential.token())) {
request = request.json(&body);
}
if let Some(key) = command.idempotency_key() {
request = request.header("Idempotency-Key", key);
}
request
}
async fn execute_watch<W: std::io::Write>(
env: &CliEnvironment,
target: WatchTarget,
options: &WatchOptions,
json: bool,
output: &mut W,
) -> Result<(), RuntimeError> {
if !json {
return Err(RuntimeError::Unavailable {
message: "watch streams JSON for agents",
next: "run logbrew watch --json",
});
}
let mut reconnect_backoff = WatchReconnectBackoff::default();
loop {
let ticket = match request_feed_ticket(env).await {
Ok(ticket) => ticket,
Err(error) if reconnect_backoff.connected_once() && !runtime_error_is_auth(&error) => {
tokio::time::sleep(reconnect_backoff.next_delay()).await;
continue;
}
Err(error) => return Err(error),
};
let live_url = feed_live_url(env.base_url.as_str(), ticket.as_str())?;
let (mut websocket, _) = match connect_async(live_url.as_str()).await {
Ok(connection) => connection,
Err(error)
if reconnect_backoff.connected_once() && !websocket_error_is_auth(&error) =>
{
tokio::time::sleep(reconnect_backoff.next_delay()).await;
continue;
}
Err(error) => return Err(map_websocket_connect_error(error)),
};
reconnect_backoff.mark_connected();
let mut emitted_before_disconnect = false;
loop {
let Some(message) = websocket.next().await else {
break;
};
let message = match message {
Ok(message) => message,
Err(error) if websocket_error_is_auth(&error) => {
return Err(map_websocket_stream_error(error));
}
Err(_) => break,
};
match message {
Message::Text(text) => {
let event = parse_live_event(text.as_str())?;
if watch_event_matches(target, options, &event) {
writeln!(output, "{event}")?;
}
emitted_before_disconnect = true;
}
Message::Binary(_) | Message::Ping(_) | Message::Pong(_) | Message::Frame(_) => {}
Message::Close(_) => return Ok(()),
}
}
if emitted_before_disconnect {
reconnect_backoff.reset();
}
tokio::time::sleep(reconnect_backoff.next_delay()).await;
}
}
#[derive(Debug, Default)]
struct WatchReconnectBackoff {
connected_once: bool,
attempts: u32,
}
impl WatchReconnectBackoff {
const fn connected_once(&self) -> bool {
self.connected_once
}
const fn mark_connected(&mut self) {
self.connected_once = true;
}
const fn reset(&mut self) {
self.attempts = 0;
}
fn next_delay(&mut self) -> std::time::Duration {
let exponent = self.attempts.min(5);
let multiplier = 1_u64 << exponent;
self.attempts = self.attempts.saturating_add(1);
let base = WATCH_RECONNECT_INITIAL_DELAY
.as_secs()
.saturating_mul(multiplier)
.min(WATCH_RECONNECT_MAX_DELAY.as_secs());
let delay = std::time::Duration::from_secs(base) + watch_reconnect_jitter();
if delay > WATCH_RECONNECT_MAX_DELAY {
WATCH_RECONNECT_MAX_DELAY
} else {
delay
}
}
}
fn watch_reconnect_jitter() -> std::time::Duration {
let Ok(elapsed) = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH) else {
return std::time::Duration::ZERO;
};
std::time::Duration::from_millis(
u64::from(elapsed.subsec_millis()) % WATCH_RECONNECT_JITTER_MAX_MILLIS,
)
}
const fn runtime_error_is_auth(error: &RuntimeError) -> bool {
matches!(
error,
RuntimeError::MissingToken | RuntimeError::Api { status: 401, .. }
)
}
fn websocket_error_is_auth(error: &WebSocketError) -> bool {
matches!(error, WebSocketError::Http(response) if response.status().as_u16() == 401)
}
async fn request_feed_ticket(env: &CliEnvironment) -> Result<String, RuntimeError> {
let url = format!("{}/api/feed/ticket", env.base_url.trim_end_matches('/'));
let client = reqwest::Client::builder()
.timeout(std::time::Duration::from_secs(30))
.connect_timeout(std::time::Duration::from_secs(10))
.build()?;
let (response, credential) =
send_authenticated_with_refresh(&client, env, |client, credential| {
client.post(url.as_str()).bearer_auth(credential.token())
})
.await?;
let status = response.status();
let body = response.text().await?;
if !status.is_success() {
return Err(RuntimeError::Api {
status: status.as_u16(),
body: credential.redact_response_body(body.as_str()),
auth_source: credential.source(),
auth_label: credential.label(),
});
}
let value = serde_json::from_str::<serde_json::Value>(body.as_str()).map_err(|_| {
RuntimeError::Unavailable {
message: "feed ticket response was not valid JSON",
next: "retry logbrew watch or run logbrew status",
}
})?;
value
.get("ticket")
.and_then(serde_json::Value::as_str)
.map(str::trim)
.filter(|ticket| !ticket.is_empty())
.map(ToOwned::to_owned)
.ok_or(RuntimeError::Unavailable {
message: "feed ticket response did not include a ticket",
next: "retry logbrew watch or run logbrew status",
})
}
fn feed_live_url(base_url: &str, ticket: &str) -> Result<String, RuntimeError> {
let trimmed = base_url.trim_end_matches('/');
let (scheme, rest) = websocket_base_parts(trimmed).ok_or(RuntimeError::Unavailable {
message: "LOGBREW_API_URL must start with http:// or https://",
next: "check LOGBREW_API_URL or run logbrew status",
})?;
Ok(format!(
"{scheme}://{rest}/api/feed/live?ticket={}",
encode_component(ticket)
))
}
fn websocket_base_parts(base_url: &str) -> Option<(&'static str, &str)> {
base_url
.strip_prefix("https://")
.map(|rest| ("wss", rest))
.or_else(|| base_url.strip_prefix("http://").map(|rest| ("ws", rest)))
}
fn parse_live_event(text: &str) -> Result<serde_json::Value, RuntimeError> {
serde_json::from_str::<serde_json::Value>(text).map_err(|_| RuntimeError::Unavailable {
message: "live watch event was not valid JSON",
next: "retry logbrew watch or check LOGBREW_API_URL",
})
}
fn watch_event_matches(
target: WatchTarget,
options: &WatchOptions,
event: &serde_json::Value,
) -> bool {
target_matches_event(target, event) && severity_matches(options, event)
}
fn target_matches_event(target: WatchTarget, event: &serde_json::Value) -> bool {
let event_type = event
.get("type")
.and_then(serde_json::Value::as_str)
.unwrap_or_default();
match target {
WatchTarget::All => true,
WatchTarget::Logs => event_type == "native_log",
WatchTarget::Issues => event_type == "native_issue",
WatchTarget::Actions => event_type == "native_action",
}
}
fn severity_matches(options: &WatchOptions, event: &serde_json::Value) -> bool {
if options.severity.is_empty() {
return true;
}
let Some(severity) = event
.get("data")
.and_then(|data| data.get("severity").or_else(|| data.get("level")))
.and_then(serde_json::Value::as_str)
else {
return false;
};
options
.severity
.iter()
.any(|allowed| allowed.as_str() == severity)
}
fn map_websocket_connect_error(error: WebSocketError) -> RuntimeError {
match error {
WebSocketError::Http(response) if response.status().as_u16() == 401 => {
RuntimeError::Unavailable {
message: "live watch ticket was rejected",
next: "run logbrew login",
}
}
WebSocketError::Http(_) => RuntimeError::Unavailable {
message: "live watch websocket upgrade failed",
next: "retry logbrew watch or check LOGBREW_API_URL",
},
WebSocketError::ConnectionClosed
| WebSocketError::AlreadyClosed
| WebSocketError::Io(_)
| WebSocketError::Tls(_)
| WebSocketError::Capacity(_)
| WebSocketError::Protocol(_)
| WebSocketError::WriteBufferFull(_)
| WebSocketError::Utf8(_)
| WebSocketError::AttackAttempt
| WebSocketError::Url(_)
| WebSocketError::HttpFormat(_) => RuntimeError::Unavailable {
message: "live watch websocket failed",
next: "retry logbrew watch or check LOGBREW_API_URL",
},
}
}
fn map_websocket_stream_error(error: WebSocketError) -> RuntimeError {
match error {
WebSocketError::ConnectionClosed | WebSocketError::AlreadyClosed => {
RuntimeError::Unavailable {
message: "live watch websocket closed",
next: "retry logbrew watch",
}
}
WebSocketError::Http(response) if response.status().as_u16() == 401 => {
RuntimeError::Unavailable {
message: "live watch ticket was rejected",
next: "run logbrew login",
}
}
WebSocketError::Http(_)
| WebSocketError::Io(_)
| WebSocketError::Tls(_)
| WebSocketError::Capacity(_)
| WebSocketError::Protocol(_)
| WebSocketError::WriteBufferFull(_)
| WebSocketError::Utf8(_)
| WebSocketError::AttackAttempt
| WebSocketError::Url(_)
| WebSocketError::HttpFormat(_) => RuntimeError::Unavailable {
message: "live watch websocket failed",
next: "retry logbrew watch or check LOGBREW_API_URL",
},
}
}
struct ReadPathFilters<'a> {
name: Option<&'a str>,
service: Option<&'a str>,
since: Option<&'a str>,
user: Option<&'a str>,
trace: Option<&'a str>,
level: Option<&'a str>,
search: Option<&'a str>,
project: Option<&'a str>,
release: Option<&'a str>,
environment: Option<&'a str>,
status: Option<&'a str>,
limit: Option<&'a str>,
min_duration_ms: Option<&'a str>,
pagination: Option<&'a str>,
cursor_time: Option<&'a str>,
cursor_id: Option<&'a str>,
}
fn read_path(target: &ReadTarget, filters: &ReadPathFilters<'_>) -> String {
match target {
ReadTarget::Logs => path_with_query(
"/api/logs",
&[
("service_name", filters.service),
("severity", filters.level),
("search", filters.search),
("since", filters.since),
("trace_id", filters.trace),
("project_id", filters.project),
("release", filters.release),
("environment", filters.environment),
("pagination", filters.pagination),
("cursor_time", filters.cursor_time),
("cursor_id", filters.cursor_id),
("limit", filters.limit),
],
),
ReadTarget::Issues => path_with_query(
"/api/telemetry/issues",
&[
("service_name", filters.service),
("since", filters.since),
("status", filters.status),
("project_id", filters.project),
("release", filters.release),
("environment", filters.environment),
("pagination", filters.pagination),
("cursor_time", filters.cursor_time),
("cursor_id", filters.cursor_id),
("limit", filters.limit),
],
),
ReadTarget::Actions => path_with_query(
"/api/telemetry/actions",
&[
("service_name", filters.service),
("name", filters.name),
("since", filters.since),
("distinct_id", filters.user),
("project_id", filters.project),
("release", filters.release),
("environment", filters.environment),
("pagination", filters.pagination),
("cursor_time", filters.cursor_time),
("cursor_id", filters.cursor_id),
("limit", filters.limit),
],
),
ReadTarget::Releases => path_with_query(
"/api/telemetry/releases",
&[
("service_name", filters.service),
("since", filters.since),
("project_id", filters.project),
("release", filters.release),
("environment", filters.environment),
("limit", filters.limit),
],
),
ReadTarget::Traces => path_with_query(
"/api/telemetry/traces",
&[
("project_id", filters.project),
("service_name", filters.service),
("release", filters.release),
("environment", filters.environment),
("status", filters.status),
("since", filters.since),
("min_duration_ms", filters.min_duration_ms),
("limit", filters.limit),
],
),
ReadTarget::Trace(id) => path_with_query(
&format!("/api/telemetry/traces/{}", encode_component(id)),
&[
("project_id", filters.project),
("release", filters.release),
("environment", filters.environment),
],
),
ReadTarget::Issue(id) => format!("/api/telemetry/issues/{}", encode_component(id)),
}
}
fn explain_path(target: &ExplainTarget) -> String {
match target {
ExplainTarget::Issue(id) => format!("/api/telemetry/issues/{}", encode_component(id)),
ExplainTarget::Trace(id) => format!("/api/telemetry/traces/{}", encode_component(id)),
}
}
fn set_path(target: &SetTarget) -> String {
match target {
SetTarget::IssueStatus { id, .. } => {
format!("/api/telemetry/issues/{}", encode_component(id))
}
}
}
fn path_with_query(path: &str, params: &[(&str, Option<&str>)]) -> String {
let query = params
.iter()
.filter_map(|(name, value)| value.map(|v| format!("{name}={}", encode_component(v))))
.collect::<Vec<_>>();
if query.is_empty() {
path.to_owned()
} else {
format!("{path}?{}", query.join("&"))
}
}
fn encode_component(value: &str) -> String {
let mut encoded = String::new();
for byte in value.bytes() {
if byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.' | b'~') {
encoded.push(char::from(byte));
} else {
encoded.push('%');
encoded.push(hex_digit(byte >> 4));
encoded.push(hex_digit(byte & 0x0f));
}
}
encoded
}
fn hex_digit(nibble: u8) -> char {
match nibble {
0..=9 => char::from(b'0' + nibble),
10..=15 => char::from(b'A' + (nibble - 10)),
_ => '?',
}
}