use anyhow::Result;
use crate::core::{
create_processing_context, init_command, set_terminal_title, set_terminal_title_and_flush,
NO_REPOS_MESSAGE, GIT_CONCURRENT_CAP,
};
use crate::git::{has_uncommitted_changes, create_and_push_tag, get_repo_visibility, RepoVisibility};
use crate::package::{get_package_info, publish_package, PublishStatus};
const SCANNING_MESSAGE: &str = "🔍 Scanning for packages...";
const PUBLISHING_MESSAGE: &str = "publishing...";
pub async fn handle_publish_command(
target_repos: Vec<String>,
dry_run: bool,
tag: bool,
allow_dirty: bool,
all: bool,
_public_only: bool,
private_only: bool,
) -> Result<()> {
set_terminal_title("📦 repos");
let (start_time, mut repos) = init_command(SCANNING_MESSAGE);
if repos.is_empty() {
println!("\r{}", NO_REPOS_MESSAGE);
set_terminal_title_and_flush("✅ repos");
return Ok(());
}
if !target_repos.is_empty() {
repos.retain(|(name, _)| {
target_repos.iter().any(|target| name.contains(target))
});
}
let filter_visibility = if all {
None } else if private_only {
Some(RepoVisibility::Private)
} else {
Some(RepoVisibility::Public)
};
use futures::stream::{FuturesUnordered, StreamExt};
use crate::package::detect_package_manager_async;
let analysis_futures: FuturesUnordered<_> = repos
.into_iter()
.map(|(name, path)| async move {
let (visibility, manager, is_dirty) = tokio::join!(
get_repo_visibility(&path),
detect_package_manager_async(&path),
async {
if !allow_dirty && !dry_run {
has_uncommitted_changes(&path).await
} else {
false
}
}
);
(name, path, visibility, manager, is_dirty)
})
.collect();
let mut analysis_results: Vec<_> = analysis_futures.collect().await;
let mut skipped_count = 0;
let mut unknown_count = 0;
if let Some(desired_visibility) = filter_visibility {
analysis_results.retain(|(_, _, visibility, _, _)| {
if *visibility == desired_visibility {
true
} else if *visibility == RepoVisibility::Unknown {
if desired_visibility == RepoVisibility::Private {
true
} else {
skipped_count += 1;
unknown_count += 1;
false
}
} else {
skipped_count += 1;
false
}
});
if skipped_count > 0 {
let visibility_type = match desired_visibility {
RepoVisibility::Public => "public",
RepoVisibility::Private => "private",
_ => "unknown",
};
let mut skip_msg = format!(
"\r📦 Filtered to {} repos only ({} skipped)",
visibility_type, skipped_count
);
if unknown_count > 0 {
skip_msg.push_str(&format!(" [{} unknown visibility treated as private]", unknown_count));
}
println!("{}\n", skip_msg);
}
}
let mut packages_to_publish = Vec::new();
let mut dirty_repos = Vec::new();
for (name, path, _visibility, manager, is_dirty) in analysis_results {
if let Some(mgr) = manager {
if is_dirty {
dirty_repos.push(name.clone());
}
packages_to_publish.push((name, path, mgr));
}
}
if !dirty_repos.is_empty() && !allow_dirty && !dry_run {
println!("\r❌ Cannot publish: {} {} uncommitted changes\n",
dirty_repos.len(),
if dirty_repos.len() == 1 { "repository has" } else { "repositories have" }
);
println!("Repositories with uncommitted changes:");
for repo in &dirty_repos {
println!(" • {}", repo);
}
println!("\nCommit your changes first, or use --allow-dirty to publish anyway (not recommended).\n");
set_terminal_title_and_flush("✅ repos");
return Ok(());
}
if packages_to_publish.is_empty() {
if target_repos.is_empty() {
println!("\r📦 No packages found in any repository\n");
} else {
println!("\r📦 No packages found matching: {}\n", target_repos.join(", "));
}
set_terminal_title_and_flush("✅ repos");
return Ok(());
}
if dry_run {
println!("\r📦 Found {} packages (dry-run mode)\n", packages_to_publish.len());
for (name, path, manager) in &packages_to_publish {
if let Some(info) = get_package_info(path).await {
println!(
" {} {:<30} ({:<7}) v{}",
manager.icon(),
name,
manager.name(),
info.version
);
} else {
println!(
" {} {:<30} ({:<7}) version unknown",
manager.icon(),
name,
manager.name()
);
}
}
println!("\nWould publish {} packages (dry-run - nothing published)\n", packages_to_publish.len());
set_terminal_title_and_flush("✅ repos");
return Ok(());
}
let total_packages = packages_to_publish.len();
let package_word = if total_packages == 1 {
"package"
} else {
"packages"
};
print!(
"\r📦 Publishing {} {} \n",
total_packages, package_word
);
println!();
let repos_for_context: Vec<(String, std::path::PathBuf)> = packages_to_publish
.iter()
.map(|(name, path, _)| (name.clone(), path.clone()))
.collect();
let context = match create_processing_context(repos_for_context, start_time, GIT_CONCURRENT_CAP) {
Ok(context) => context,
Err(e) => {
set_terminal_title_and_flush("✅ repos");
return Err(e);
}
};
process_publish_repositories(context, packages_to_publish, tag).await;
set_terminal_title_and_flush("✅ repos");
Ok(())
}
#[derive(Default)]
struct PublishStatistics {
published: usize,
already_published: usize,
errors: Vec<(String, String)>,
}
impl PublishStatistics {
fn update(&mut self, status: &PublishStatus, repo_name: &str, message: &str) {
match status {
PublishStatus::Published => self.published += 1,
PublishStatus::AlreadyPublished => self.already_published += 1,
PublishStatus::Error => self.errors.push((repo_name.to_string(), message.to_string())),
_ => {}
}
}
fn generate_summary(&self, _total: usize) -> String {
let mut parts = Vec::new();
if self.published > 0 {
parts.push(format!("✅ {} published", self.published));
}
if self.already_published > 0 {
parts.push(format!("⚠️ {} already published", self.already_published));
}
if !self.errors.is_empty() {
parts.push(format!("❌ {} failed", self.errors.len()));
}
if parts.is_empty() {
"No changes".to_string()
} else {
parts.join(" ")
}
}
}
async fn process_publish_repositories(
context: crate::core::ProcessingContext,
packages: Vec<(String, std::path::PathBuf, crate::package::PackageManager)>,
tag: bool,
) {
use crate::core::{create_progress_bar};
use futures::stream::{FuturesUnordered, StreamExt};
use std::sync::{Arc, Mutex};
let mut futures = FuturesUnordered::new();
let statistics = Arc::new(Mutex::new(PublishStatistics::default()));
let mut repo_progress_bars = Vec::new();
for (repo_name, _) in &context.repositories {
let progress_bar =
create_progress_bar(&context.multi_progress, &context.progress_style, repo_name);
progress_bar.set_message(PUBLISHING_MESSAGE);
repo_progress_bars.push(progress_bar);
}
let _separator_pb = crate::core::create_separator_progress_bar(&context.multi_progress);
let footer_pb = crate::core::create_footer_progress_bar(&context.multi_progress);
footer_pb.set_message("Starting...");
let _separator_pb2 = crate::core::create_separator_progress_bar(&context.multi_progress);
let max_name_length = context.max_name_length;
let total_packages = packages.len();
let publish_semaphore = Arc::new(tokio::sync::Semaphore::new(8));
for (((repo_name, repo_path, manager), progress_bar), _) in
packages.into_iter().zip(repo_progress_bars).zip(context.repositories)
{
let stats_clone = Arc::clone(&statistics);
let semaphore_clone = Arc::clone(&publish_semaphore);
let footer_clone = footer_pb.clone();
let future = async move {
let _permit = semaphore_clone.acquire().await.unwrap();
let (success, message) = publish_package(&repo_path, &manager, false).await;
let status = if success {
if message.contains("already") {
PublishStatus::AlreadyPublished
} else {
PublishStatus::Published
}
} else {
PublishStatus::Error
};
let mut final_message = message.clone();
if tag && matches!(status, PublishStatus::Published) {
if let Some(info) = get_package_info(&repo_path).await {
let tag_name = format!("v{}", info.version);
let (tag_success, tag_message) = create_and_push_tag(&repo_path, &tag_name).await;
if tag_success {
final_message = format!("{}, {}", message, tag_message);
} else {
final_message = format!("{} (tag failed: {})", message, tag_message);
}
}
}
progress_bar.set_prefix(format!(
"{} {:width$}",
status.symbol(),
repo_name,
width = max_name_length
));
progress_bar.set_message(format!("{:<20} {}", status.text(), final_message));
progress_bar.finish();
{
let mut stats_guard = stats_clone.lock().unwrap();
stats_guard.update(&status, &repo_name, &final_message);
let summary = stats_guard.generate_summary(total_packages);
footer_clone.set_message(summary);
}
};
futures.push(future);
}
while futures.next().await.is_some() {}
footer_pb.finish();
let final_stats = statistics.lock().unwrap();
if !final_stats.errors.is_empty() {
println!("\n{}", "━".repeat(70));
println!("❌ Failed to publish:\n");
for (repo, error) in &final_stats.errors {
println!(" • {}: {}", repo, error);
}
println!("{}", "━".repeat(70));
}
println!();
}