From b8eb2279f211ec4b5a8d6bac9b92d8b6c19fcb8f Mon Sep 17 00:00:00 2001 From: dzejkop Date: Thu, 9 Jul 2026 09:48:09 +0200 Subject: [PATCH] Add progress indicators for long-running CLI actions --- crates/tangled-cli/src/commands/spindle.rs | 14 +-- crates/tangled-cli/src/main.rs | 11 +- crates/tangled-cli/src/ops/auth.rs | 15 ++- crates/tangled-cli/src/ops/mod.rs | 27 +++++ crates/tangled-cli/src/ops/oauth.rs | 123 +++++++++++++-------- crates/tangled-cli/src/ops/pull.rs | 2 +- crates/tangled-cli/src/ops/repo.rs | 57 ++++++++-- crates/tangled-cli/src/ops/secrets.rs | 7 +- crates/tangled-cli/src/ops/session.rs | 5 +- crates/tangled-cli/src/progress.rs | 53 +++++++++ 10 files changed, 234 insertions(+), 80 deletions(-) create mode 100644 crates/tangled-cli/src/progress.rs diff --git a/crates/tangled-cli/src/commands/spindle.rs b/crates/tangled-cli/src/commands/spindle.rs index 7445728..9325e4c 100644 --- a/crates/tangled-cli/src/commands/spindle.rs +++ b/crates/tangled-cli/src/commands/spindle.rs @@ -166,15 +166,13 @@ async fn logs(args: SpindleLogsArgs) -> Result<()> { let ws_url = format!("{}/spindle/logs/{}/{}/{}", spindle_base, knot, rkey, name); - println!( - "Connecting to logs stream for {}:{}:{}...", - knot, rkey, name - ); - // Connect to WebSocket - let (ws_stream, _) = connect_async(&ws_url) - .await - .map_err(|e| anyhow!("Failed to connect to log stream: {}", e))?; + let (ws_stream, _) = crate::progress::with_spinner( + format!("Connecting to logs stream for {knot}:{rkey}:{name}..."), + connect_async(&ws_url), + ) + .await + .map_err(|e| anyhow!("Failed to connect to log stream: {}", e))?; let (mut _write, mut read) = ws_stream.split(); diff --git a/crates/tangled-cli/src/main.rs b/crates/tangled-cli/src/main.rs index f9c34ac..ffaf5f7 100644 --- a/crates/tangled-cli/src/main.rs +++ b/crates/tangled-cli/src/main.rs @@ -1,6 +1,7 @@ mod cli; mod commands; mod ops; +mod progress; mod util; use clap::Parser; @@ -13,8 +14,12 @@ async fn main() -> Result<()> { human_panic::setup_panic!(); let cli = Cli::parse(); - commands::dispatch(cli) - .await - .map_err(|err| color_eyre::eyre::eyre!("{err:?}"))?; + if let Err(err) = commands::dispatch(cli).await { + eprintln!("Error: {err}"); + for cause in err.chain().skip(1) { + eprintln!(" caused by: {cause}"); + } + std::process::exit(1); + } Ok(()) } diff --git a/crates/tangled-cli/src/ops/auth.rs b/crates/tangled-cli/src/ops/auth.rs index 750e4b2..a1d8c9d 100644 --- a/crates/tangled-cli/src/ops/auth.rs +++ b/crates/tangled-cli/src/ops/auth.rs @@ -45,18 +45,27 @@ impl PdsAuth { &self, rb: reqwest::RequestBuilder, ) -> Result { - Ok(xrpc::parse(self.response(rb).await?).await?) + crate::progress::with_spinner("Contacting PDS...", async { + Ok(xrpc::parse(self.response(rb).await?).await?) + }) + .await } pub async fn send_unit(&self, rb: reqwest::RequestBuilder) -> Result<()> { - Ok(xrpc::parse_unit(self.response(rb).await?).await?) + crate::progress::with_spinner("Contacting PDS...", async { + Ok(xrpc::parse_unit(self.response(rb).await?).await?) + }) + .await } pub async fn send_bytes( &self, rb: reqwest::RequestBuilder, ) -> Result> { - Ok(xrpc::parse_bytes(self.response(rb).await?).await?) + crate::progress::with_spinner("Contacting PDS...", async { + Ok(xrpc::parse_bytes(self.response(rb).await?).await?) + }) + .await } async fn response( diff --git a/crates/tangled-cli/src/ops/mod.rs b/crates/tangled-cli/src/ops/mod.rs index de26bea..c91da44 100644 --- a/crates/tangled-cli/src/ops/mod.rs +++ b/crates/tangled-cli/src/ops/mod.rs @@ -17,6 +17,9 @@ pub mod types; use std::sync::OnceLock; +use serde::de::DeserializeOwned; +use tangled_api::xrpc::{self, XrpcError}; + /// Base URL for Tangled's read-only XRPC API/AppView (Bobbin). pub const DEFAULT_API_BASE: &str = "https://api.tangled.org"; @@ -37,6 +40,30 @@ pub(crate) fn http() -> &'static reqwest::Client { }) } +pub(crate) async fn xrpc_send( + rb: reqwest::RequestBuilder, +) -> Result +where + T: DeserializeOwned, +{ + crate::progress::with_spinner("Contacting server...", xrpc::send(rb)).await +} + +pub(crate) async fn xrpc_send_unit( + rb: reqwest::RequestBuilder, +) -> Result<(), XrpcError> { + crate::progress::with_spinner("Contacting server...", xrpc::send_unit(rb)) + .await +} + +#[allow(dead_code)] +pub(crate) async fn xrpc_send_bytes( + rb: reqwest::RequestBuilder, +) -> Result, XrpcError> { + crate::progress::with_spinner("Contacting server...", xrpc::send_bytes(rb)) + .await +} + /// Extracts the rkey from an at-uri (`at://did/collection/rkey`). pub(crate) fn uri_rkey(uri: &str) -> Option { uri.rsplit('/').next().map(|s| s.to_string()) diff --git a/crates/tangled-cli/src/ops/oauth.rs b/crates/tangled-cli/src/ops/oauth.rs index b9e5017..96d9fd5 100644 --- a/crates/tangled-cli/src/ops/oauth.rs +++ b/crates/tangled-cli/src/ops/oauth.rs @@ -48,17 +48,20 @@ async fn resolve_identity(input: &str) -> Result<(String, String)> { struct Res { did: String, } - let res = http() - .get(format!( - "https://public.api.bsky.app/xrpc/com.atproto.identity.resolveHandle?handle={}", - urlencode(input) - )) - .send() - .await?; - if !res.status().is_success() { - bail!("could not resolve handle {input}: {}", res.status()); - } - res.json::().await?.did + crate::progress::with_spinner("Resolving handle...", async { + let res = http() + .get(format!( + "https://public.api.bsky.app/xrpc/com.atproto.identity.resolveHandle?handle={}", + urlencode(input) + )) + .send() + .await?; + if !res.status().is_success() { + bail!("could not resolve handle {input}: {}", res.status()); + } + Ok::<_, anyhow::Error>(res.json::().await?.did) + }) + .await? }; let doc_url = if let Some(plc) = did.strip_prefix("did:plc:") { @@ -69,7 +72,12 @@ async fn resolve_identity(input: &str) -> Result<(String, String)> { bail!("unsupported DID method: {did}"); }; let doc: serde_json::Value = - http().get(&doc_url).send().await?.json().await?; + crate::progress::with_spinner("Fetching DID document...", async { + Ok::<_, anyhow::Error>( + http().get(&doc_url).send().await?.json().await?, + ) + }) + .await?; let pds = doc["service"] .as_array() .and_then(|services| { @@ -99,30 +107,44 @@ async fn discover_auth_server(pds: &str) -> Result { struct ProtectedResource { authorization_servers: Vec, } - let protected_resource: ProtectedResource = http() - .get(format!( - "{}/.well-known/oauth-protected-resource", - pds.trim_end_matches('/') - )) - .send() - .await? - .json() - .await - .context("fetching PDS protected-resource metadata")?; + let protected_resource: ProtectedResource = + crate::progress::with_spinner("Fetching OAuth metadata...", async { + Ok::<_, anyhow::Error>( + http() + .get(format!( + "{}/.well-known/oauth-protected-resource", + pds.trim_end_matches('/') + )) + .send() + .await? + .json() + .await + .context("fetching PDS protected-resource metadata")?, + ) + }) + .await?; let issuer = protected_resource .authorization_servers .first() .ok_or_else(|| anyhow!("PDS lists no authorization servers"))?; - let meta: AuthServerMetadata = http() - .get(format!( - "{}/.well-known/oauth-authorization-server", - issuer.trim_end_matches('/') - )) - .send() - .await? - .json() - .await - .context("fetching authorization server metadata")?; + let meta: AuthServerMetadata = crate::progress::with_spinner( + "Fetching authorization server metadata...", + async { + Ok::<_, anyhow::Error>( + http() + .get(format!( + "{}/.well-known/oauth-authorization-server", + issuer.trim_end_matches('/') + )) + .send() + .await? + .json() + .await + .context("fetching authorization server metadata")?, + ) + }, + ) + .await?; Ok(meta) } @@ -136,21 +158,30 @@ async fn as_form_post( ) -> Result { for _ in 0..2 { let proof = key.proof("POST", url, nonce.as_deref(), None); - let res = http() - .post(url) - .header("DPoP", proof) - .form(form) - .send() - .await?; - if let Some(n) = res - .headers() - .get("dpop-nonce") - .and_then(|v| v.to_str().ok()) - { - *nonce = Some(n.to_string()); + let (status, body, new_nonce) = crate::progress::with_spinner( + "Contacting authorization server...", + async { + let res = http() + .post(url) + .header("DPoP", proof) + .form(form) + .send() + .await?; + let new_nonce = res + .headers() + .get("dpop-nonce") + .and_then(|v| v.to_str().ok()) + .map(str::to_string); + let status = res.status(); + let body: serde_json::Value = + res.json().await.unwrap_or_default(); + Ok::<_, reqwest::Error>((status, body, new_nonce)) + }, + ) + .await?; + if let Some(n) = new_nonce { + *nonce = Some(n); } - let status = res.status(); - let body: serde_json::Value = res.json().await.unwrap_or_default(); if status.is_success() { return Ok(body); } diff --git a/crates/tangled-cli/src/ops/pull.rs b/crates/tangled-cli/src/ops/pull.rs index 4fecdd9..9a1998b 100644 --- a/crates/tangled-cli/src/ops/pull.rs +++ b/crates/tangled-cli/src/ops/pull.rs @@ -263,7 +263,7 @@ pub async fn merge_pull( author_email: None, }; let rb = tangled_repo::merge(http(), knot_base, &input).bearer_auth(&token); - let _: serde_json::Value = tangled_api::xrpc::send(rb).await?; + let _: serde_json::Value = super::xrpc_send(rb).await?; Ok(()) } diff --git a/crates/tangled-cli/src/ops/repo.rs b/crates/tangled-cli/src/ops/repo.rs index aacd538..609e2dc 100644 --- a/crates/tangled-cli/src/ops/repo.rs +++ b/crates/tangled-cli/src/ops/repo.rs @@ -113,12 +113,25 @@ pub async fn create_repo( swap_commit: None, }; let rb = atproto_repo::create_record(http(), opts.pds_base, &input); - let created: atproto_repo::create_record::Output = - opts.auth.send(rb).await?; - - // Extract rkey from at-uri: at://did/collection/rkey - let rkey = uri_rkey(&created.uri) - .ok_or_else(|| anyhow!("failed to parse rkey from uri"))?; + let created: Result = + opts.auth.send(rb).await; + + // Extract rkey from at-uri: at://did/collection/rkey. Some PDSes + // occasionally return a 500 after committing createRecord; because we use + // a deterministic rkey, verify the record before failing hard. + let rkey = match created { + Ok(created) => uri_rkey(&created.uri) + .ok_or_else(|| anyhow!("failed to parse rkey from uri"))?, + Err(err) => { + if repo_record_exists(opts.pds_base, opts.did, opts.name, opts.auth) + .await + { + opts.name.to_string() + } else { + return Err(err); + } + } + }; // 2) Obtain a service auth token for the knot (aud = did:web:) let token = service_auth::mint( @@ -139,7 +152,8 @@ pub async fn create_repo( }; let rb = tangled_repo::create(http(), knot_base, &input).bearer_auth(&token); - let created_repo: tangled_repo::create::Output = xrpc::send(rb).await?; + let created_repo: tangled_repo::create::Output = + super::xrpc_send(rb).await?; if let Some(repo_did) = created_repo.repo_did.as_deref() { update_repo_record(opts.pds_base, opts.did, &rkey, opts.auth, |obj| { obj.insert("repoDid".to_string(), serde_json::json!(repo_did)); @@ -149,6 +163,24 @@ pub async fn create_repo( Ok(()) } +async fn repo_record_exists( + pds_base: &str, + did: &str, + rkey: &str, + auth: &PdsAuth, +) -> bool { + let params = atproto_repo::get_record::Params { + repo: did.to_string(), + collection: REPO_COLLECTION.to_string(), + rkey: rkey.to_string(), + cid: None, + }; + let rb = atproto_repo::get_record(http(), pds_base, ¶ms); + auth.send::(rb) + .await + .is_ok() +} + /// Fetches the sh.tangled.repo record, applies `mutate` to it, and writes it /// back with putRecord. Unknown fields are preserved because the record is /// round-tripped as raw JSON. @@ -290,7 +322,7 @@ async fn describe_repo_via_knot( if let Some(token) = bearer { rb = rb.bearer_auth(token); } - Ok(xrpc::send(rb).await?) + Ok(super::xrpc_send(rb).await?) } fn is_default_bobbin_base(base: &str) -> bool { @@ -320,7 +352,7 @@ async fn get_repo_by_repo_did( let rb = http() .get(xrpc::xrpc_url(api_base, "sh.tangled.repo.getRepoByRepoDid")) .query(&[("repoDid", repo_did)]); - let out: BobbinRepoRecord = xrpc::send(rb).await?; + let out: BobbinRepoRecord = super::xrpc_send(rb).await?; let owner_did = uri_did(&out.uri) .ok_or_else(|| anyhow!("Bobbin repo response uri missing owner DID"))?; let rkey = uri_rkey(&out.uri) @@ -365,7 +397,7 @@ pub async fn delete_repo( }; let rb = tangled_repo::delete(http(), knot_base, &input).bearer_auth(&token); - xrpc::send_unit(rb).await?; + super::xrpc_send_unit(rb).await?; Ok(()) } @@ -408,7 +440,8 @@ pub async fn get_default_branch( repo: format!("{}/{}", did, name), }; let rb = tangled_repo::get_default_branch(http(), knot_host, ¶ms); - let out: tangled_repo::get_default_branch::Output = xrpc::send(rb).await?; + let out: tangled_repo::get_default_branch::Output = + super::xrpc_send(rb).await?; Ok(DefaultBranch { name: out.name, hash: out.hash, @@ -429,7 +462,7 @@ pub async fn get_languages( }; let rb = tangled_repo::languages(http(), knot_host, ¶ms); // Live knots don't always return the full lexicon shape; stay lenient. - let res: serde_json::Value = xrpc::send(rb).await?; + let res: serde_json::Value = super::xrpc_send(rb).await?; let langs = res .get("languages") .cloned() diff --git a/crates/tangled-cli/src/ops/secrets.rs b/crates/tangled-cli/src/ops/secrets.rs index 021c9a2..484cf8b 100644 --- a/crates/tangled-cli/src/ops/secrets.rs +++ b/crates/tangled-cli/src/ops/secrets.rs @@ -2,7 +2,6 @@ use anyhow::Result; use tangled_api::lexicon::sh::tangled::repo as tangled_repo; -use tangled_api::xrpc; use super::auth::PdsAuth; use super::types::Secret; @@ -26,7 +25,7 @@ pub async fn list_repo_secrets( }; let rb = tangled_repo::list_secrets(http(), service_base, ¶ms) .bearer_auth(&token); - let res: tangled_repo::list_secrets::Output = xrpc::send(rb).await?; + let res: tangled_repo::list_secrets::Output = super::xrpc_send(rb).await?; Ok(res.secrets) } @@ -52,7 +51,7 @@ pub async fn add_repo_secret( }; let rb = tangled_repo::add_secret(http(), service_base, &input) .bearer_auth(&token); - xrpc::send_unit(rb).await?; + super::xrpc_send_unit(rb).await?; Ok(()) } @@ -76,6 +75,6 @@ pub async fn remove_repo_secret( }; let rb = tangled_repo::remove_secret(http(), service_base, &input) .bearer_auth(&token); - xrpc::send_unit(rb).await?; + super::xrpc_send_unit(rb).await?; Ok(()) } diff --git a/crates/tangled-cli/src/ops/session.rs b/crates/tangled-cli/src/ops/session.rs index 53727a0..15617da 100644 --- a/crates/tangled-cli/src/ops/session.rs +++ b/crates/tangled-cli/src/ops/session.rs @@ -2,7 +2,6 @@ use anyhow::Result; use tangled_api::lexicon::com::atproto::server; -use tangled_api::xrpc; use tangled_config::session::Session; pub async fn login_with_password( @@ -16,7 +15,7 @@ pub async fn login_with_password( ..Default::default() }; let rb = server::create_session(super::http(), pds_base, &input); - let out: server::create_session::Output = xrpc::send(rb).await?; + let out: server::create_session::Output = super::xrpc_send(rb).await?; Ok(Session { access_jwt: out.access_jwt, refresh_jwt: out.refresh_jwt, @@ -32,7 +31,7 @@ pub async fn refresh_session( ) -> Result { let rb = server::refresh_session(super::http(), pds_base) .bearer_auth(refresh_jwt); - let out: server::refresh_session::Output = xrpc::send(rb).await?; + let out: server::refresh_session::Output = super::xrpc_send(rb).await?; Ok(Session { access_jwt: out.access_jwt, refresh_jwt: out.refresh_jwt, diff --git a/crates/tangled-cli/src/progress.rs b/crates/tangled-cli/src/progress.rs new file mode 100644 index 0000000..7e8b7d7 --- /dev/null +++ b/crates/tangled-cli/src/progress.rs @@ -0,0 +1,53 @@ +use std::future::Future; +use std::io::{self, IsTerminal}; +use std::time::Duration; + +use indicatif::{ProgressBar, ProgressStyle}; + +/// Runs `future` while displaying a transient spinner on stderr. +/// +/// Progress indicators are hidden automatically when stderr is not a TTY, so +/// JSON output, tests, and redirected command output stay deterministic. +pub async fn with_spinner( + message: impl Into, + future: impl Future, +) -> T { + let pb = spinner(message); + let result = future.await; + pb.finish_and_clear(); + result +} + +pub fn spinner(message: impl Into) -> ProgressBar { + let pb = if io::stderr().is_terminal() { + ProgressBar::new_spinner() + } else { + ProgressBar::hidden() + }; + pb.set_style( + ProgressStyle::with_template("{spinner:.cyan} {msg}") + .expect("spinner template is valid") + .tick_strings(&["⠋", "⠙", "⠹", "⠸", "⠼", "⠴", "⠦", "⠧", "⠇", "⠏"]), + ); + pb.set_message(message.into()); + pb.enable_steady_tick(Duration::from_millis(100)); + pb +} + +pub fn progress_bar(message: impl Into) -> ProgressBar { + let pb = if io::stderr().is_terminal() { + ProgressBar::new(0) + } else { + ProgressBar::hidden() + }; + pb.set_style( + ProgressStyle::with_template( + "{spinner:.cyan} {msg} [{wide_bar:.cyan/blue}] {pos}/{len}", + ) + .expect("progress bar template is valid") + .progress_chars("=>-"), + ); + pb.set_message(message.into()); + pb.enable_steady_tick(Duration::from_millis(100)); + pb +} -- 2.51.2