MCPcopy Create free account
hub / github.com/brocode/fw / synchronize

Function synchronize

src/sync/mod.rs:61–133  ·  view source on GitHub ↗
(maybe_config: Result<Config, AppError>, only_new: bool, ff_merge: bool, tags: &BTreeSet<String>, worker: i32)

Source from the content-addressed store, hash-verified

59}
60
61pub 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 }));

Callers 1

_mainFunction · 0.85

Calls 2

ssh_agent_runningFunction · 0.85
sync_projectFunction · 0.85

Tested by

no test coverage detected