1use super::*;
2
3async fn discover_packages(
16 workspace_manifest: &Path,
17) -> Result<(Vec<Package>, bool), PublishError> {
18 let content: String = read_to_string(workspace_manifest).await?;
19 let doc: Value = toml::from_str(&content)
20 .map_err(|_error: toml::de::Error| PublishError::ManifestParseError)?;
21 let workspace_version: Option<String> = doc
22 .get(TOML_WORKSPACE)
23 .and_then(|workspace: &Value| workspace.get(TOML_PACKAGE))
24 .and_then(|package: &Value| package.get(TOML_VERSION))
25 .and_then(|version: &Value| version.as_str())
26 .map(|version: &str| version.to_string());
27 let mut packages: Vec<Package> = Vec::new();
28 if let Some(workspace) = doc.get(TOML_WORKSPACE)
29 && let Some(members) = workspace
30 .get(TOML_MEMBERS)
31 .and_then(|members_value: &Value| members_value.as_array())
32 {
33 for member in members {
34 if let Some(pattern) = member.as_str() {
35 let base_path: &Path = workspace_manifest.parent().unwrap_or(workspace_manifest);
36 expand_pattern(
37 base_path,
38 pattern,
39 &mut packages,
40 workspace_version.as_deref(),
41 )
42 .await?;
43 }
44 }
45 }
46 let has_root_package: bool = doc.get(TOML_PACKAGE).is_some();
47 if has_root_package {
48 let root_package: Package =
49 read_package_manifest(workspace_manifest, workspace_version.as_deref()).await?;
50 packages.push(root_package);
51 }
52 Ok((packages, has_root_package))
53}
54
55async fn expand_pattern(
68 base_path: &Path,
69 pattern: &str,
70 packages: &mut Vec<Package>,
71 workspace_version: Option<&str>,
72) -> Result<(), PublishError> {
73 if pattern.contains('*') {
74 let parent: &Path = Path::new(pattern).parent().unwrap_or(Path::new("."));
75 let full_parent: PathBuf = base_path.join(parent);
76 if full_parent.is_dir() {
77 let mut entries: ReadDir = read_dir(&full_parent).await?;
78 while let Some(entry) = entries.next_entry().await? {
79 let path: PathBuf = entry.path();
80 if path.is_dir() {
81 let cargo_toml: PathBuf = path.join(CARGO_TOML);
82 if cargo_toml.exists() {
83 let package: Package =
84 read_package_manifest(&cargo_toml, workspace_version).await?;
85 packages.push(package);
86 }
87 }
88 }
89 }
90 } else {
91 let cargo_toml: PathBuf = base_path.join(pattern).join(CARGO_TOML);
92 if cargo_toml.exists() {
93 let package: Package = read_package_manifest(&cargo_toml, workspace_version).await?;
94 packages.push(package);
95 }
96 }
97 Ok(())
98}
99
100async fn read_package_manifest(
114 manifest_path: &Path,
115 workspace_version: Option<&str>,
116) -> Result<Package, PublishError> {
117 let content: String = read_to_string(manifest_path).await?;
118 let doc: Value = toml::from_str(&content)
119 .map_err(|_error: toml::de::Error| PublishError::ManifestParseError)?;
120 let package_table: &Value = doc
121 .get(TOML_PACKAGE)
122 .ok_or(PublishError::ManifestParseError)?;
123 let name: String = package_table
124 .get(TOML_NAME)
125 .and_then(|n: &Value| n.as_str())
126 .ok_or(PublishError::ManifestParseError)?
127 .to_string();
128 let version: String = match package_table.get(TOML_VERSION) {
129 Some(version_value) => {
130 if let Some(version_str) = version_value.as_str() {
131 version_str.to_string()
132 } else if version_value
133 .get(TOML_WORKSPACE)
134 .and_then(|workspace_value: &Value| workspace_value.as_bool())
135 .unwrap_or(false)
136 {
137 workspace_version
138 .ok_or(PublishError::ManifestParseError)?
139 .to_string()
140 } else {
141 return Err(PublishError::ManifestParseError);
142 }
143 }
144 None => workspace_version
145 .ok_or(PublishError::ManifestParseError)?
146 .to_string(),
147 };
148 let publish: bool = package_table
149 .get(TOML_PUBLISH_KEY)
150 .and_then(|publish_value: &Value| publish_value.as_bool())
151 .unwrap_or(true);
152 let path: PathBuf = manifest_path
153 .parent()
154 .filter(|p: &&Path| !p.as_os_str().is_empty())
155 .map_or_else(|| PathBuf::from("."), |p: &Path| p.to_path_buf());
156 let local_dependencies: Vec<String> = extract_local_dependencies(&doc, manifest_path)?;
157 Ok(Package {
158 name,
159 version,
160 path,
161 local_dependencies,
162 publish,
163 })
164}
165
166fn extract_local_dependencies(
183 doc: &Value,
184 _manifest_path: &Path,
185) -> Result<Vec<String>, PublishError> {
186 let mut deps: Vec<String> = Vec::new();
187 let dep_sections: [&str; 3] = [
188 TOML_DEPENDENCIES,
189 TOML_BUILD_DEPENDENCIES,
190 TOML_DEV_DEPENDENCIES,
191 ];
192 for section in &dep_sections {
193 if let Some(table) = doc
194 .get(section)
195 .and_then(|section_value: &Value| section_value.as_table())
196 {
197 for (dep_name, dep_value) in table {
198 let is_local: bool = match dep_value {
199 Value::Table(t) => {
200 let has_path_or_workspace: bool = t.get(TOML_PATH).is_some()
201 || t.get(TOML_WORKSPACE)
202 .and_then(|workspace_value: &Value| workspace_value.as_bool())
203 .unwrap_or(false);
204 let versioned: bool = t.get(TOML_VERSION).is_some();
205 has_path_or_workspace && (*section != TOML_DEV_DEPENDENCIES || versioned)
206 }
207 _ => false,
208 };
209 if is_local {
210 deps.push(dep_name.clone());
211 }
212 }
213 }
214 }
215 Ok(deps)
216}
217
218fn validate_publish_order(packages: &[Package]) -> Result<(), PublishError> {
231 let position: HashMap<String, usize> = packages
232 .iter()
233 .enumerate()
234 .map(|(index, package): (usize, &Package)| (package.name.clone(), index))
235 .collect();
236 for package in packages {
237 let Some(package_position) = position.get(&package.name) else {
238 continue;
239 };
240 for dep in &package.local_dependencies {
241 if let Some(dep_position) = position.get(dep)
242 && dep_position > package_position
243 {
244 return Err(PublishError::InvalidPublishOrder(format!(
245 "{} depends on {} but is listed before it in [workspace.members]",
246 package.name, dep
247 )));
248 }
249 }
250 }
251 Ok(())
252}
253
254fn position_root_package(packages: &mut Vec<Package>) -> Result<(), PublishError> {
273 let Some(root) = packages.pop() else {
274 return Ok(());
275 };
276 let earliest_dependent: Option<usize> = packages
277 .iter()
278 .enumerate()
279 .filter(|(_index, package): &(usize, &Package)| {
280 package.local_dependencies.contains(&root.name)
281 })
282 .map(|(index, _item): (usize, &Package)| index)
283 .min();
284 let Some(earliest) = earliest_dependent else {
285 packages.push(root);
286 return Ok(());
287 };
288 let member_positions: HashMap<&str, usize> = packages
289 .iter()
290 .enumerate()
291 .map(|(index, package): (usize, &Package)| (package.name.as_str(), index))
292 .collect();
293 if let Some(max_dep) = root
294 .local_dependencies
295 .iter()
296 .filter_map(|dep: &String| member_positions.get(dep.as_str()))
297 .max()
298 && max_dep >= &earliest
299 {
300 return Err(PublishError::InvalidPublishOrder(format!(
301 "{} must publish after its dependency at members position {} but before dependent at position {}; reorder [workspace.members]",
302 root.name, max_dep, earliest
303 )));
304 }
305 packages.insert(earliest, root);
306 Ok(())
307}
308
309pub async fn resolve_publish_order(manifest_path: &str) -> Result<Vec<Package>, PublishError> {
322 let workspace_manifest: &Path = Path::new(manifest_path);
323 let (mut packages, has_root_package) = discover_packages(workspace_manifest).await?;
324 if has_root_package {
325 position_root_package(&mut packages)?;
326 }
327 validate_publish_order(&packages)?;
328 Ok(packages)
329}
330
331pub fn is_already_published(stderr: &str) -> bool {
351 stderr.contains(STDERR_ALREADY_EXISTS_ON)
352 || stderr.contains(STDERR_ALREADY_BEEN_UPLOADED)
353 || stderr.contains(STDERR_IS_ALREADY_PUBLISHED)
354}
355
356pub fn is_rate_limited(stderr: &str) -> bool {
375 stderr.contains(STDERR_TOO_MANY_REQUESTS) || stderr.contains(STDERR_TOO_MANY_NEW_CRATES)
376}
377
378fn civil_to_days(year: i64, month: i64, day: i64) -> i64 {
393 let shifted_year: i64 = if month <= 2 { year - 1 } else { year };
394 let era: i64 = shifted_year.div_euclid(ERA_YEARS);
395 let year_of_era: i64 = shifted_year.rem_euclid(ERA_YEARS);
396 let month_position: i64 = (month + 9) % 12;
397 let day_of_year: i64 =
398 (MONTH_POSITION_SCALE * month_position + MONTH_POSITION_ROUNDING) / MONTH_POSITION_DIVISOR;
399 let day_of_era: i64 = year_of_era * DAYS_PER_YEAR + year_of_era.div_euclid(LEAP_CYCLE_YEARS)
400 - year_of_era.div_euclid(CENTURY_YEARS)
401 + day_of_year
402 + day
403 - 1;
404 era * ERA_DAYS + day_of_era - CIVIL_EPOCH_OFFSET_DAYS
405}
406
407pub fn parse_rate_limit_wait_secs(stderr: &str) -> Option<u64> {
422 let start: usize = stderr.find(STDERR_TRY_AGAIN_AFTER)? + STDERR_TRY_AGAIN_AFTER.len();
423 let end: usize = start + stderr[start..].find(STDERR_TRY_AGAIN_AFTER_END)?;
424 let timestamp: &str = stderr[start..end].trim();
425 let mut tokens: std::str::SplitWhitespace<'_> = timestamp.split_whitespace();
426 let _: &str = tokens.next()?;
427 let day: i64 = tokens
428 .next()?
429 .trim_end_matches(DAY_FIELD_SUFFIX)
430 .parse()
431 .ok()?;
432 let month_token: &str = tokens.next()?;
433 let month: i64 = MONTH_ABBREVIATIONS
434 .iter()
435 .position(|name: &&str| *name == month_token)? as i64
436 + 1;
437 let year: i64 = tokens.next()?.parse().ok()?;
438 let mut clock_fields: std::str::Split<'_, char> = tokens.next()?.split(CLOCK_FIELD_SEPARATOR);
439 let hour: i64 = clock_fields.next()?.parse().ok()?;
440 let minute: i64 = clock_fields.next()?.parse().ok()?;
441 let second: i64 = clock_fields.next()?.parse().ok()?;
442 let day_start: i64 = civil_to_days(year, month, day) * SECONDS_PER_DAY;
443 let deadline: i64 = day_start + hour * SECONDS_PER_HOUR + minute * SECONDS_PER_MINUTE + second;
444 let now: i64 = std::time::SystemTime::now()
445 .duration_since(std::time::UNIX_EPOCH)
446 .ok()?
447 .as_secs() as i64;
448 let remaining: i64 = deadline - now;
449 let bounded: i64 = remaining.clamp(0, RATE_LIMIT_MAX_WAIT_SECS as i64);
454 Some(bounded as u64 + RATE_LIMIT_SKEW_SECS)
455}
456
457fn rate_limit_wait_secs(stderr: &str) -> u64 {
467 parse_rate_limit_wait_secs(stderr).unwrap_or(RATE_LIMIT_FLOOR_SECS + RATE_LIMIT_SKEW_SECS)
468}
469
470async fn publish_package_with_retry(package: &Package, max_retries: u32) -> PublishResult {
489 let mut attempt: u32 = 0;
490 let mut last_error: String;
491 let mut rate_limited_attempts: u32 = 0;
492 loop {
493 match publish_single_package(package).await {
494 Ok(()) => {
495 return PublishResult {
496 package_name: package.name.clone(),
497 success: true,
498 error: None,
499 retries: attempt + rate_limited_attempts,
500 };
501 }
502 Err(error) => {
503 let stderr: String = error.to_string();
504 if is_rate_limited(&stderr) && rate_limited_attempts < RATE_LIMIT_MAX_WAITS {
505 let wait: u64 = rate_limit_wait_secs(&stderr);
506 log::info!(
507 "{}: rate limited by the registry, waiting {}s for the deadline it reported",
508 package.name,
509 wait
510 );
511 rate_limited_attempts += 1;
512 sleep(Duration::from_secs(wait)).await;
513 continue;
514 }
515 attempt += 1;
516 last_error = stderr;
517 if attempt <= max_retries {
518 sleep(Duration::from_secs(2_u64.pow(attempt))).await;
519 continue;
520 }
521 break;
522 }
523 }
524 }
525 PublishResult {
526 package_name: package.name.clone(),
527 success: false,
528 error: Some(last_error),
529 retries: attempt + rate_limited_attempts,
530 }
531}
532
533async fn publish_single_package(package: &Package) -> Result<(), Box<dyn std::error::Error>> {
543 let output: std::process::Output = Command::new(CARGO)
544 .arg(CARGO_PUBLISH)
545 .arg(CLI_FLAG_ALLOW_DIRTY)
546 .arg(CLI_FLAG_NO_VERIFY)
547 .current_dir(&package.path)
548 .stdout(Stdio::piped())
549 .stderr(Stdio::piped())
550 .output()
551 .await?;
552 if output.status.success() {
553 return Ok(());
554 }
555 let stderr: String = String::from_utf8_lossy(&output.stderr).to_string();
556 if is_already_published(&stderr) {
557 log::info!("{} is already published, treating as success", package.name);
558 return Ok(());
559 }
560 Err(stderr.into())
561}
562
563pub async fn execute_publish(
578 manifest_path: &str,
579 max_retries: u32,
580) -> Result<Vec<PublishResult>, PublishError> {
581 let path: &Path = Path::new(manifest_path);
582 let path: &Path = match path.parent() {
583 Some(parent) if !parent.as_os_str().is_empty() => parent,
584 _ => Path::new("."),
585 };
586 let workspace_manifest: PathBuf = path.join(CARGO_TOML);
587 let sync_report: SyncReport =
588 match execute_sync(workspace_manifest.to_str().unwrap_or(CARGO_TOML)).await {
589 Ok(report) => report,
590 Err(error) => return Err(PublishError::SyncFailed(error)),
591 };
592 if sync_report.file_changed {
593 log::info!(
594 "publish: synced workspace dependencies ({} renamed, {} versioned) to v{}",
595 sync_report.renamed_entries.len(),
596 sync_report.versioned_entries.len(),
597 sync_report.workspace_version,
598 );
599 }
600 let ordered_packages: Vec<Package> =
601 resolve_publish_order(workspace_manifest.to_str().unwrap_or(CARGO_TOML)).await?;
602 if ordered_packages.is_empty() {
603 return Ok(Vec::new());
604 }
605 let mut results: Vec<PublishResult> = Vec::new();
606 for package in ordered_packages {
607 if !package.publish {
608 log::info!("Skipping {} (publish = false)", package.name);
609 continue;
610 }
611 log::info!("Publishing {} v{}...", package.name, package.version);
612 let result: PublishResult = publish_package_with_retry(&package, max_retries).await;
613 if result.success {
614 if result.retries == 0 {
615 log::info!("Successfully published {}", result.package_name,);
616 } else {
617 log::info!(
618 "Successfully published {} (retried {} times)",
619 result.package_name,
620 result.retries
621 );
622 }
623 } else if let Some(error) = &result.error {
624 log::error!("Failed to publish {}: {error}", result.package_name);
625 } else {
626 log::error!("Failed to publish {}", result.package_name);
627 }
628 results.push(result);
629 }
630 Ok(results)
631}