use anyhow::Result;
use crate::core::{
create_processing_context, init_command, set_terminal_title, set_terminal_title_and_flush,
NO_REPOS_MESSAGE,
};
use crate::git::Status;
const SCANNING_MESSAGE: &str = "🔍 Scanning for git repositories...";
pub async fn handle_push_command(
force_push: bool,
verbose: bool,
show_changes: bool,
no_drift_check: bool,
jobs: Option<usize>,
sequential: bool,
) -> Result<()> {
use crate::core::config::get_git_concurrency;
set_terminal_title("🚀 repos");
let (start_time, repos) = init_command(SCANNING_MESSAGE);
if repos.is_empty() {
println!("\r{}", NO_REPOS_MESSAGE);
set_terminal_title_and_flush("✅ repos");
return Ok(());
}
let concurrent_limit = get_git_concurrency(jobs, sequential);
let total_repos = repos.len();
let repo_word = if total_repos == 1 {
"repository"
} else {
"repositories"
};
let concurrency_info = if verbose {
format!(" ({} concurrent)", concurrent_limit)
} else {
String::new()
};
print!(
"\r🚀 Pushing {} {}{} \n",
total_repos, repo_word, concurrency_info
);
println!();
let context = match create_processing_context(repos, start_time, concurrent_limit) {
Ok(context) => context,
Err(e) => {
set_terminal_title_and_flush("✅ repos");
return Err(e);
}
};
process_push_repositories(context, force_push, verbose, show_changes).await;
if !no_drift_check {
check_and_display_drift();
}
set_terminal_title_and_flush("✅ repos");
Ok(())
}
async fn process_push_repositories(context: crate::core::ProcessingContext, force_push: bool, verbose: bool, show_changes: bool) {
use crate::core::{acquire_stats_lock, create_progress_bar};
use crate::git::{fetch_and_analyze, push_if_needed};
use futures::stream::{FuturesUnordered, StreamExt};
let repo_progress_bars: Vec<_> = if verbose {
context.repositories.iter()
.map(|(repo_name, _)| {
let pb = create_progress_bar(&context.multi_progress, &context.progress_style, repo_name);
pb.set_message("processing...");
pb
})
.collect()
} else {
use indicatif::{ProgressBar, ProgressStyle};
let single_pb = context.multi_progress.add(ProgressBar::new(context.total_repos as u64));
single_pb.set_style(ProgressStyle::default_bar().template("[{pos}/{len}] {msg}").unwrap());
single_pb.set_message("📤 Processing...");
vec![single_pb; context.repositories.len()]
};
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("✅ 0 Pushed 🟢 0 Synced 🔴 0 Failed 🟡 0 No Upstream 🟠 0 Skipped".to_string());
let _separator_pb2 = crate::core::create_separator_progress_bar(&context.multi_progress);
let max_name_length = context.max_name_length;
let start_time = context.start_time;
let total_repos = context.total_repos;
use crate::core::config::FETCH_CONCURRENT_CAP;
let fetch_concurrency = (context.max_concurrency * 2).min(FETCH_CONCURRENT_CAP);
let fetch_semaphore = std::sync::Arc::new(tokio::sync::Semaphore::new(fetch_concurrency));
let rate_limit_count = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let has_rate_limit = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let mut pipeline_futures = FuturesUnordered::new();
for ((repo_name, repo_path), progress_bar) in context.repositories.into_iter().zip(repo_progress_bars) {
let fetch_semaphore_clone = std::sync::Arc::clone(&fetch_semaphore);
let push_semaphore_clone = std::sync::Arc::clone(&context.semaphore);
let stats_clone = std::sync::Arc::clone(&context.statistics);
let footer_clone = footer_pb.clone();
let rate_limit_count_clone = std::sync::Arc::clone(&rate_limit_count);
let has_rate_limit_clone = std::sync::Arc::clone(&has_rate_limit);
let verbose_clone = verbose;
let max_name_length_clone = max_name_length;
let start_time_clone = start_time;
let total_repos_clone = total_repos;
let future = async move {
use crate::core::config::SLOW_REPO_THRESHOLD_SECS;
let repo_start_time = std::time::Instant::now();
let _fetch_permit = match fetch_semaphore_clone.acquire().await {
Ok(permit) => permit,
Err(e) => {
eprintln!("Error: Failed to acquire fetch permit for {}: {}", repo_name, e);
let stats = acquire_stats_lock(&stats_clone);
stats.update(
&repo_name,
&repo_path.to_string_lossy(),
&Status::Error,
&format!("semaphore error: {}", e),
false,
);
if verbose_clone {
progress_bar.finish_with_message(format!("🔴 {} semaphore error", repo_name));
}
return;
}
};
let fetch_result = fetch_and_analyze(&repo_path, force_push).await;
drop(_fetch_permit);
let _push_permit = match push_semaphore_clone.acquire().await {
Ok(permit) => permit,
Err(e) => {
eprintln!("Error: Failed to acquire push permit for {}: {}", repo_name, e);
let stats = acquire_stats_lock(&stats_clone);
stats.update(
&repo_name,
&repo_path.to_string_lossy(),
&Status::Error,
&format!("semaphore error: {}", e),
false,
);
if verbose_clone {
progress_bar.finish_with_message(format!("🔴 {} semaphore error", repo_name));
}
return;
}
};
let mut attempt = 0;
let max_attempts = 2;
let result = loop {
attempt += 1;
let (status, message, has_uncommitted) = push_if_needed(&repo_path, &fetch_result, force_push).await;
if message.contains("⚠️ RATE LIMIT") {
has_rate_limit_clone.store(true, std::sync::atomic::Ordering::Release);
rate_limit_count_clone.fetch_add(1, std::sync::atomic::Ordering::Release);
if attempt < max_attempts {
tokio::time::sleep(tokio::time::Duration::from_secs(2)).await;
continue;
} else {
let suggestion = format!(
"{} (try reducing concurrency with --jobs N or --sequential)",
message.replace("⚠️ RATE LIMIT: ", "")
);
break (status, suggestion, has_uncommitted);
}
}
break (status, message, has_uncommitted);
};
let (status, message, has_uncommitted_changes) = result;
let repo_elapsed = repo_start_time.elapsed();
let repo_elapsed_secs = repo_elapsed.as_secs_f32();
let display_message = if has_uncommitted_changes && matches!(status, crate::git::Status::Synced) {
format!("{} (uncommitted changes)", message)
} else {
message.clone()
};
let display_message = if repo_elapsed.as_secs() >= SLOW_REPO_THRESHOLD_SECS {
format!("{} ({:.1}s)", display_message, repo_elapsed_secs)
} else {
display_message
};
if verbose_clone {
progress_bar.set_prefix(format!("{} {:width$}", status.symbol(), repo_name, width = max_name_length_clone));
progress_bar.set_message(format!("{:<10} {}", status.text(), display_message));
progress_bar.finish();
} else {
progress_bar.set_message(format!("{} {} ({})", status.symbol(), repo_name, status.text()));
progress_bar.inc(1);
}
let stats = acquire_stats_lock(&stats_clone);
stats.update(&repo_name, &repo_path.to_string_lossy(), &status, &message, has_uncommitted_changes);
let duration = start_time_clone.elapsed();
if verbose_clone {
footer_clone.set_message(stats.generate_summary(total_repos_clone, duration));
} else {
use std::sync::atomic::Ordering;
let live_counters = format!(
"✅ {} Pushed 🟢 {} Synced 🔴 {} Failed 🟡 {} No Upstream 🟠 {} Skipped",
stats.total_commits_pushed.load(Ordering::Relaxed),
stats.synced_repos.load(Ordering::Relaxed),
stats.error_repos.load(Ordering::Relaxed),
stats.no_upstream_repos.lock().unwrap().len(),
stats.skipped_repos.load(Ordering::Relaxed)
);
footer_clone.set_message(live_counters);
}
};
pipeline_futures.push(future);
}
while pipeline_futures.next().await.is_some() {}
if has_rate_limit.load(std::sync::atomic::Ordering::Acquire) {
let count = rate_limit_count.load(std::sync::atomic::Ordering::Acquire);
eprintln!("\n⚠️ Rate limit detected on {} operation(s).", count);
eprintln!("💡 Try reducing concurrency: repos push --jobs 3");
}
footer_pb.finish();
let final_stats = acquire_stats_lock(&context.statistics);
let detailed_summary = final_stats.generate_detailed_summary(show_changes);
if !detailed_summary.is_empty() {
println!("\n{}", "━".repeat(70));
println!("{}", detailed_summary);
println!("{}", "━".repeat(70));
}
println!();
}
fn check_and_display_drift() {
if let Ok(statuses) = crate::subrepo::status::analyze_subrepos() {
if statuses.iter().any(|s| s.has_drift) {
crate::subrepo::status::display_drift_summary(&statuses);
}
}
}
pub async fn handle_pull_command(
use_rebase: bool,
verbose: bool,
show_changes: bool,
no_drift_check: bool,
jobs: Option<usize>,
sequential: bool,
) -> Result<()> {
use crate::core::config::get_git_concurrency;
set_terminal_title("🔽 repos");
let (start_time, repos) = init_command(SCANNING_MESSAGE);
if repos.is_empty() {
println!("\r{}", NO_REPOS_MESSAGE);
set_terminal_title_and_flush("✅ repos");
return Ok(());
}
let concurrent_limit = get_git_concurrency(jobs, sequential);
let total_repos = repos.len();
let repo_word = if total_repos == 1 {
"repository"
} else {
"repositories"
};
let concurrency_info = if verbose {
format!(" ({} concurrent)", concurrent_limit)
} else {
String::new()
};
let pull_strategy = if use_rebase { " with rebase" } else { "" };
print!(
"\r🔽 Pulling {} {}{}{} \n",
total_repos, repo_word, pull_strategy, concurrency_info
);
println!();
let context = match create_processing_context(repos, start_time, concurrent_limit) {
Ok(context) => context,
Err(e) => {
set_terminal_title_and_flush("✅ repos");
return Err(e);
}
};
process_pull_repositories(context, use_rebase, verbose, show_changes).await;
if !no_drift_check {
check_and_display_drift();
}
set_terminal_title_and_flush("✅ repos");
Ok(())
}
async fn process_pull_repositories(context: crate::core::ProcessingContext, use_rebase: bool, verbose: bool, show_changes: bool) {
use crate::core::{acquire_stats_lock, create_progress_bar};
use crate::git::{fetch_and_analyze_for_pull, pull_if_needed};
use futures::stream::{FuturesUnordered, StreamExt};
let repo_progress_bars: Vec<_> = if verbose {
context.repositories.iter()
.map(|(repo_name, _)| {
let pb = create_progress_bar(&context.multi_progress, &context.progress_style, repo_name);
pb.set_message("processing...");
pb
})
.collect()
} else {
use indicatif::{ProgressBar, ProgressStyle};
let single_pb = context.multi_progress.add(ProgressBar::new(context.total_repos as u64));
single_pb.set_style(ProgressStyle::default_bar().template("[{pos}/{len}] {msg}").unwrap());
single_pb.set_message("🔽 Processing...");
vec![single_pb; context.repositories.len()]
};
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("🔽 0 Pulled 🟢 0 Synced 🔴 0 Failed 🟡 0 No Upstream 🟠 0 Skipped".to_string());
let _separator_pb2 = crate::core::create_separator_progress_bar(&context.multi_progress);
let max_name_length = context.max_name_length;
let start_time = context.start_time;
let total_repos = context.total_repos;
use crate::core::config::FETCH_CONCURRENT_CAP;
let fetch_concurrency = (context.max_concurrency * 2).min(FETCH_CONCURRENT_CAP);
let fetch_semaphore = std::sync::Arc::new(tokio::sync::Semaphore::new(fetch_concurrency));
let rate_limit_count = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let has_rate_limit = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let total_commits_pulled = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let mut pipeline_futures = FuturesUnordered::new();
for ((repo_name, repo_path), progress_bar) in context.repositories.into_iter().zip(repo_progress_bars) {
let fetch_semaphore_clone = std::sync::Arc::clone(&fetch_semaphore);
let pull_semaphore_clone = std::sync::Arc::clone(&context.semaphore);
let stats_clone = std::sync::Arc::clone(&context.statistics);
let footer_clone = footer_pb.clone();
let rate_limit_count_clone = std::sync::Arc::clone(&rate_limit_count);
let has_rate_limit_clone = std::sync::Arc::clone(&has_rate_limit);
let total_commits_pulled_clone = std::sync::Arc::clone(&total_commits_pulled);
let verbose_clone = verbose;
let max_name_length_clone = max_name_length;
let start_time_clone = start_time;
let total_repos_clone = total_repos;
let future = async move {
use crate::core::config::SLOW_REPO_THRESHOLD_SECS;
let repo_start_time = std::time::Instant::now();
let _fetch_permit = match fetch_semaphore_clone.acquire().await {
Ok(permit) => permit,
Err(e) => {
eprintln!("Error: Failed to acquire fetch permit for {}: {}", repo_name, e);
let stats = acquire_stats_lock(&stats_clone);
stats.update(
&repo_name,
&repo_path.to_string_lossy(),
&Status::Error,
&format!("semaphore error: {}", e),
false,
);
if verbose_clone {
progress_bar.finish_with_message(format!("🔴 {} semaphore error", repo_name));
}
return;
}
};
let fetch_result = fetch_and_analyze_for_pull(&repo_path).await;
drop(_fetch_permit);
let _pull_permit = match pull_semaphore_clone.acquire().await {
Ok(permit) => permit,
Err(e) => {
eprintln!("Error: Failed to acquire pull permit for {}: {}", repo_name, e);
let stats = acquire_stats_lock(&stats_clone);
stats.update(
&repo_name,
&repo_path.to_string_lossy(),
&Status::Error,
&format!("semaphore error: {}", e),
false,
);
if verbose_clone {
progress_bar.finish_with_message(format!("🔴 {} semaphore error", repo_name));
}
return;
}
};
let mut attempt = 0;
let max_attempts = 2;
let result = loop {
attempt += 1;
let (status, message, has_uncommitted) = pull_if_needed(&repo_path, &fetch_result, use_rebase).await;
if message.contains("⚠️ RATE LIMIT") {
has_rate_limit_clone.store(true, std::sync::atomic::Ordering::Release);
rate_limit_count_clone.fetch_add(1, std::sync::atomic::Ordering::Release);
if attempt < max_attempts {
tokio::time::sleep(tokio::time::Duration::from_secs(2)).await;
continue;
} else {
let suggestion = format!(
"{} (try reducing concurrency with --jobs N or --sequential)",
message.replace("⚠️ RATE LIMIT: ", "")
);
break (status, suggestion, has_uncommitted);
}
}
break (status, message, has_uncommitted);
};
let (status, message, has_uncommitted_changes) = result;
if matches!(status, crate::git::Status::Pulled) {
total_commits_pulled_clone.fetch_add(fetch_result.behind_count as usize, std::sync::atomic::Ordering::Relaxed);
}
let repo_elapsed = repo_start_time.elapsed();
let repo_elapsed_secs = repo_elapsed.as_secs_f32();
let display_message = if has_uncommitted_changes && matches!(status, crate::git::Status::Synced) {
format!("{} (uncommitted changes)", message)
} else {
message.clone()
};
let display_message = if repo_elapsed.as_secs() >= SLOW_REPO_THRESHOLD_SECS {
format!("{} ({:.1}s)", display_message, repo_elapsed_secs)
} else {
display_message
};
if verbose_clone {
progress_bar.set_prefix(format!("{} {:width$}", status.symbol(), repo_name, width = max_name_length_clone));
progress_bar.set_message(format!("{:<10} {}", status.text(), display_message));
progress_bar.finish();
} else {
progress_bar.set_message(format!("{} {} ({})", status.symbol(), repo_name, status.text()));
progress_bar.inc(1);
}
let stats = acquire_stats_lock(&stats_clone);
stats.update(&repo_name, &repo_path.to_string_lossy(), &status, &message, has_uncommitted_changes);
let duration = start_time_clone.elapsed();
if verbose_clone {
footer_clone.set_message(stats.generate_summary(total_repos_clone, duration));
} else {
use std::sync::atomic::Ordering;
let live_counters = format!(
"🔽 {} Pulled 🟢 {} Synced 🔴 {} Failed 🟡 {} No Upstream 🟠 {} Skipped",
total_commits_pulled_clone.load(Ordering::Relaxed),
stats.synced_repos.load(Ordering::Relaxed),
stats.error_repos.load(Ordering::Relaxed),
stats.no_upstream_repos.lock().unwrap().len(),
stats.skipped_repos.load(Ordering::Relaxed)
);
footer_clone.set_message(live_counters);
}
};
pipeline_futures.push(future);
}
while pipeline_futures.next().await.is_some() {}
if has_rate_limit.load(std::sync::atomic::Ordering::Acquire) {
let count = rate_limit_count.load(std::sync::atomic::Ordering::Acquire);
eprintln!("\n⚠️ Rate limit detected on {} operation(s).", count);
eprintln!("💡 Try reducing concurrency: repos pull --jobs 3");
}
footer_pb.finish();
let final_stats = acquire_stats_lock(&context.statistics);
let detailed_summary = final_stats.generate_detailed_summary(show_changes);
if !detailed_summary.is_empty() {
println!("\n{}", "━".repeat(70));
println!("{}", detailed_summary);
println!("{}", "━".repeat(70));
}
println!();
}