(maybe_config: Result<Config, AppError>, only_new: bool, ff_merge: bool, tags: &BTreeSet<String>, worker: i32)
| 59 | } |
| 60 | |
| 61 | pub fn synchronize(maybe_config: Result<Config, AppError>, only_new: bool, ff_merge: bool, tags: &BTreeSet<String>, worker: i32) -> Result<(), AppError> { |
| 62 | eprintln!("Synchronizing everything"); |
| 63 | if !ssh_agent_running() { |
| 64 | eprintln!("SSH Agent not running. Process may hang.") |
| 65 | } |
| 66 | let config = Arc::new(maybe_config?); |
| 67 | |
| 68 | let projects: Vec<Project> = config.projects.values().map(ToOwned::to_owned).collect(); |
| 69 | let q: Arc<SegQueue<Project>> = Arc::new(SegQueue::new()); |
| 70 | let projects_count = projects.len() as u64; |
| 71 | |
| 72 | projects |
| 73 | .into_iter() |
| 74 | .filter(|p| tags.is_empty() || p.tags.clone().unwrap_or_default().intersection(tags).count() > 0) |
| 75 | .for_each(|p| q.push(p)); |
| 76 | |
| 77 | let spinner_style = ProgressStyle::default_spinner() |
| 78 | .tick_chars("⣾⣽⣻⢿⡿⣟⣯⣷⣿") |
| 79 | .template("{prefix:.bold.dim} {spinner} {wide_msg}") |
| 80 | .map_err(|e| AppError::RuntimeError(format!("Invalid Template: {}", e)))?; |
| 81 | |
| 82 | let m = MultiProgress::new(); |
| 83 | m.set_draw_target(ProgressDrawTarget::stderr()); |
| 84 | |
| 85 | let job_results: Arc<SegQueue<Result<(), AppError>>> = Arc::new(SegQueue::new()); |
| 86 | let progress_bars = (1..=worker).map(|i| { |
| 87 | let pb = m.add(ProgressBar::new(projects_count)); |
| 88 | pb.set_style(spinner_style.clone()); |
| 89 | pb.set_prefix(format!("[{: >2}/{}]", i, worker)); |
| 90 | pb.set_message("initializing..."); |
| 91 | pb.tick(); |
| 92 | pb.enable_steady_tick(Duration::from_millis(250)); |
| 93 | pb |
| 94 | }); |
| 95 | let mut thread_handles: Vec<thread::JoinHandle<()>> = Vec::new(); |
| 96 | for pb in progress_bars { |
| 97 | let job_q = Arc::clone(&q); |
| 98 | let job_config = Arc::clone(&config); |
| 99 | let job_result_queue = Arc::clone(&job_results); |
| 100 | thread_handles.push(thread::spawn(move || { |
| 101 | let mut job_result: Result<(), AppError> = Result::Ok(()); |
| 102 | loop { |
| 103 | if let Some(project) = job_q.pop() { |
| 104 | pb.set_message(project.name.to_string()); |
| 105 | let sync_result = sync_project(&job_config, &project, only_new, ff_merge); |
| 106 | let msg = match sync_result { |
| 107 | Ok(_) => format!("DONE: {}", project.name), |
| 108 | Err(ref e) => format!("FAILED: {} - {}", project.name, e), |
| 109 | }; |
| 110 | pb.println(&msg); |
| 111 | job_result = job_result.and(sync_result); |
| 112 | } else { |
| 113 | pb.finish_and_clear(); |
| 114 | break; |
| 115 | } |
| 116 | } |
| 117 | job_result_queue.push(job_result); |
| 118 | })); |
no test coverage detected