Skip to main content

crate_cli/publish/
fn.rs

1use super::*;
2
3/// Discover all packages in the workspace: every `[workspace.members]`
4/// entry in declaration order, with the root package (when the workspace
5/// root manifest also declares `[package]`) appended last.
6///
7/// # Arguments
8///
9/// - `&Path` - Path to workspace root Cargo.toml
10///
11/// # Returns
12///
13/// - `Result<(Vec<Package>, bool), PublishError>` - Packages and whether a
14///   root package was appended
15async 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
55/// Expand glob pattern to find package directories
56///
57/// # Arguments
58///
59/// - `&Path` - Base path for expansion
60/// - `&str` - Glob pattern
61/// - `&mut Vec<Package>` - Output vector for found packages
62/// - `Option<&str>` - Workspace root version for `version.workspace = true`
63///
64/// # Returns
65///
66/// - `Result<(), PublishError>` - Success or error
67async 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
100/// Read package manifest and extract information
101///
102/// `version.workspace = true` resolves to the workspace root version
103/// passed in `workspace_version`.
104///
105/// # Arguments
106///
107/// - `&Path` - Path to package Cargo.toml
108/// - `Option<&str>` - Workspace root `[workspace.package].version`, if any
109///
110/// # Returns
111///
112/// - `Result<Package, PublishError>` - Package info or error
113async 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
166/// Extract local workspace dependencies that constrain publish order
167///
168/// `[dependencies]` and `[build-dependencies]` entries with `path` or
169/// `workspace = true` always constrain. `[dev-dependencies]` are stripped
170/// from the published manifest, so they constrain only when they carry a
171/// `version` field (cargo publish registry-checks versioned dev-deps);
172/// path-only dev-deps skip the registry and impose no order constraint.
173///
174/// # Arguments
175///
176/// - `&Value` - Parsed manifest
177/// - `&Path` - Path to manifest for resolving relative paths
178///
179/// # Returns
180///
181/// - `Result<Vec<String>, PublishError>` - List of local dependency names
182fn 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
218/// Validate that the publish order satisfies every package's local
219/// dependency constraints: a package must never appear before a
220/// workspace-local dependency of its own.
221///
222/// # Arguments
223///
224/// - `&[Package]` - Packages in intended publish order
225///
226/// # Returns
227///
228/// - `Result<(), PublishError>` - `InvalidPublishOrder` naming the first
229///   offending pair when the order violates a local dependency.
230fn 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
254/// Move the root package (appended last by `discover_packages`) to its
255/// topological position when workspace members depend on it
256///
257/// When no member depends on the root package the root stays last (the
258/// conventional facade-last layout). When members do depend on the root
259/// (e.g. `ui` / `engine` crates depending on a root facade crate), the
260/// root is inserted right before the earliest such member, provided all
261/// of the root's own local dependencies appear earlier in the members
262/// order; otherwise the members order cannot satisfy both constraints
263/// and `InvalidPublishOrder` is returned.
264///
265/// # Arguments
266///
267/// - `&mut Vec<Package>` - Packages with the root package as last element
268///
269/// # Returns
270///
271/// - `Result<(), PublishError>` - Success or `InvalidPublishOrder`
272fn 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
309/// Resolve the publish order for a workspace: `[workspace.members]`
310/// declaration order, with the root package (if any) placed at its
311/// topological position, validated against local dependency constraints.
312///
313/// # Arguments
314///
315/// - `&str` - Path to the workspace root Cargo.toml
316///
317/// # Returns
318///
319/// - `Result<Vec<Package>, PublishError>` - Ordered packages, or an
320///   error when the members order violates a local dependency.
321pub 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
331/// Check whether `cargo publish` stderr indicates the package version is
332/// already present on the registry (a success case for idempotent
333/// re-runs).
334///
335/// A version that is already live is not a failure: the artifact the caller
336/// asked for is on the registry, so the release has achieved what it came
337/// for. Republishing a live version is refused by the registry rather than
338/// allowed to overwrite it, which makes a re-run fail on the 24 crates that
339/// shipped before the rate limit refused the remaining three. Treating the
340/// refusal as success is what lets a re-run converge on the whole workspace
341/// instead of stalling on work already done.
342///
343/// # Arguments
344///
345/// - `&str` - cargo publish stderr output
346///
347/// # Returns
348///
349/// - `bool` - True when the output reports the version is already published
350pub 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
356/// Check whether `cargo publish` stderr reports a registry rate-limit
357/// refusal rather than a fault that retrying cannot fix.
358///
359/// crates.io meters **new** crate names on its own schedule: a small burst
360/// allowance, then roughly one name per ten minutes. Once a workspace
361/// publishes more new names than the burst allows, every further name is
362/// refused with a deadline embedded in the message. That deadline is
363/// already longer than any exponential backoff, so a retry that ignores it
364/// only spends the budget and fails again — this is the class of failure
365/// that must wait out the deadline the registry handed back.
366///
367/// # Arguments
368///
369/// - `&str` - cargo publish stderr output
370///
371/// # Returns
372///
373/// - `bool` - True when the output reports a registry rate limit
374pub fn is_rate_limited(stderr: &str) -> bool {
375    stderr.contains(STDERR_TOO_MANY_REQUESTS) || stderr.contains(STDERR_TOO_MANY_NEW_CRATES)
376}
377
378/// Convert a civil date to a day count since 1970-01-01.
379///
380/// Counts from 0000-03-01 so that a leap day lands at the end of a year
381/// and the month positions become a fixed-length stride.
382///
383/// # Arguments
384///
385/// - `i64` - Proleptic Gregorian year
386/// - `i64` - Month, 1 through 12
387/// - `i64` - Day of month
388///
389/// # Returns
390///
391/// - `i64` - Days since the Unix epoch
392fn 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
407/// Parse the RFC 2822 retry deadline crates.io embeds in a refusal.
408///
409/// The message reads `Please try again after Tue, 29 Sep 2026 04:55:57 GMT`,
410/// where the day-of-week prefix and the `GMT` suffix are noise. A timestamp
411/// that does not match that shape yields `None` so the caller falls back to
412/// the conservative floor instead of guessing a shorter wait.
413///
414/// # Arguments
415///
416/// - `&str` - cargo publish stderr output
417///
418/// # Returns
419///
420/// - `Option<u64>` - Seconds to wait, or `None` when no deadline is present
421pub 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    // A parsed deadline is honoured as given, in either direction: a window
450    // the registry says is shorter than the skew allowance is already open
451    // on arrival, and a guess would only extend a wait it priced precisely.
452    // The floor is for the unreadable message, handled by the caller.
453    let bounded: i64 = remaining.clamp(0, RATE_LIMIT_MAX_WAIT_SECS as i64);
454    Some(bounded as u64 + RATE_LIMIT_SKEW_SECS)
455}
456
457/// Seconds to wait before retrying a refused publish.
458///
459/// # Arguments
460///
461/// - `&str` - cargo publish stderr output
462///
463/// # Returns
464///
465/// - `u64` - Wait bounded by the registry deadline when it carried one
466fn 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
470/// Publish a single package with retry logic
471///
472/// A registry rate limit is not a fault: the registry named a deadline, and
473/// retrying before it only spends the budget to be refused again. Those
474/// refusals wait out the deadline the registry itself reported and are not
475/// charged against `max_retries`, because the budget exists to survive
476/// transient errors — a metered limit is neither transient nor shorter than
477/// any backoff this loop could choose. Every other failure keeps the original
478/// exponential backoff and still counts.
479///
480/// # Arguments
481///
482/// - `&Package` - Package to publish
483/// - `u32` - Maximum retry attempts
484///
485/// # Returns
486///
487/// - `PublishResult` - Result with success status and retry count
488async 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
533/// Execute cargo publish command for a single package
534///
535/// # Arguments
536///
537/// - `&Package` - Package to publish
538///
539/// # Returns
540///
541/// - `Result<(), Box<dyn std::error::Error>>` - Success or error
542async 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
563/// Execute publish command for all packages in workspace
564///
565/// Publishes in `[workspace.members]` declaration order with the root
566/// package (if any) last, after validating the order against local
567/// dependency constraints.
568///
569/// # Arguments
570///
571/// - `&str` - Path to workspace Cargo.toml
572/// - `u32` - Maximum retry attempts per package
573///
574/// # Returns
575///
576/// - `Result<Vec<PublishResult>, PublishError>` - Results for all packages
577pub 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}