Skip to main content

mkit_cli/remote_dispatch/
mod.rs

1//! URL-scheme → `Transport` dispatch for `mkit push` / `mkit pull`.
2//!
3//! The Rust binary wires all five shipping schemes here: `mkit+file://`,
4//! `mkit+https://` (and `mkit+http://` for local dev), `mkit+s3://`, and
5//! `mkit+ssh://`. The memory transport is in-process only, so it is
6//! reached via [`push_all`] / [`pull_all`] with an `Arc<MemoryTransport>`
7//! constructed in-process rather than URL-based construction. Integration
8//! tests in the `mkit-cli` crate exercise the memory path directly.
9//!
10//! Credentials / environment sources:
11//! - HTTP(S): optional `MKIT_API_TOKEN` bearer.
12//! - S3/R2: `MKIT_R2_ACCESS_KEY_ID` + `MKIT_R2_SECRET_ACCESS_KEY` (plus
13//!   optional `MKIT_R2_REGION`, default `auto`). Missing creds do NOT
14//!   fail at connect time; the first signed request returns
15//!   `TransportError::AccessDenied`.
16//! - SSH: spawns `ssh(1)` subprocess — inherits the user's agent / keys /
17//!   `~/.ssh/config`. Per-repo `.mkit/config` SSH options (host-key
18//!   checking, known-hosts path, identity file) are wired through via
19//!   `SshTransport::connect_with_options` when config is loaded.
20
21// `pub(crate)` so the `remote remove`/`rename` command handlers can drive
22// the record's lifecycle ops (#545); everything else stays module-private.
23pub(crate) mod applied_packs;
24mod envelope_signer;
25pub(crate) mod grants;
26mod packmap;
27mod split;
28mod upload_receipts;
29
30use mkit_core::layout::RepoLayout;
31use std::path::Path;
32use std::sync::Arc;
33
34use applied_packs::AppliedPacks;
35
36use mkit_core::hash::{HASH_LEN, Hash};
37use mkit_core::object::Object;
38use mkit_core::ops::merge::is_ancestor;
39use mkit_core::ops::restore;
40use mkit_core::pack::{self, PackError, PackWriter, PreparedDelta, PreparedRaw};
41use mkit_core::protocol::{PackKey, Transport, TransportError, UploadLimits};
42use mkit_core::refs::{self, Head};
43use mkit_core::store::{ObjectStore, StoreError};
44use mkit_core::transfer::{self, PackListError};
45use mkit_transport_connect::admission::{AdmissionPolicy, is_reserved};
46use mkit_transport_connect::{
47    ConnectTransport, PENDING_INTERRUPTED_MESSAGE, repository_identity_from_url,
48};
49use mkit_transport_file::FileTransport;
50use mkit_transport_s3::S3Transport;
51use mkit_transport_ssh::{SshInitError, SshOptions, SshTransport, parse_mkit_ssh_url};
52use rayon::prelude::*;
53
54pub use split::{MAX_CHAIN_COMMITS, MAX_SPLIT_STEPS, PushControl, StepAuthority, plan_push_steps};
55
56use packmap::{
57    ChainAction, advance_packmap, apply_fetched_chain, commit_head, packmap_ref, probe_chain,
58    rebaseline_depth, resolve_and_download_chain,
59};
60
61const DEFAULT_REMOTE: &str = "default";
62
63/// Errors returned by the push / pull helpers. Mapped to exit codes by
64/// the commands themselves.
65#[derive(Debug, thiserror::Error)]
66pub enum DispatchError {
67    #[error("unsupported URL scheme: {0}")]
68    UnsupportedScheme(String),
69    #[error("malformed URL: {0}")]
70    MalformedUrl(String),
71    #[error(
72        "repository `{identity}` not found at {origin}; check the remote URL path against the server's configured name (an empty path selects `default` on a single-repository deployment)"
73    )]
74    RepositoryNotFound { identity: String, origin: String },
75    #[error("no HEAD branch to push")]
76    NoHead,
77    /// A poll-loop checkpoint observed `signal::is_shutdown() == true`
78    /// and aborted partway through. Callers should map this to
79    /// `exit::TEMPFAIL` (75) so retries are safe — the transfer is
80    /// half-finished but the remote is unmodified for any ref we
81    /// hadn't reached yet.
82    #[error("interrupted")]
83    Interrupted,
84    #[error("{0}")]
85    UploadInterrupted(String),
86    #[error("transport: {0}")]
87    Transport(#[from] TransportError),
88    #[error("refs: {0}")]
89    Refs(#[from] refs::RefError),
90    #[error("repo lock: {0}")]
91    RepoLock(#[from] mkit_core::repo_lock::LockError),
92    #[error("worktree discovery: {0}")]
93    Discover(#[from] mkit_core::layout::DiscoverError),
94    #[error("io: {0}")]
95    Io(#[from] std::io::Error),
96    #[error("store: {0}")]
97    Store(#[from] StoreError),
98    #[error("pack: {0}")]
99    Pack(#[from] PackError),
100    #[error("packlist: {0}")]
101    PackList(#[from] PackListError),
102    #[error("ssh init: {0}")]
103    SshInit(#[from] SshInitError),
104    #[error("pull requires HEAD to point at a branch")]
105    DetachedHead,
106    #[error("remote branch '{0}' not found")]
107    RemoteBranchMissing(String),
108    #[error("pull would not fast-forward branch '{branch}'; merge or rebase first")]
109    NonFastForwardPull { branch: String },
110    #[error("restore safety: {0}")]
111    RestoreSafety(String),
112    #[error("object is not a commit")]
113    NotCommit,
114    #[error("restore: {0}")]
115    Restore(#[from] restore::RestoreError),
116    /// The per-endpoint credential-trust gate (#97) refused to build a
117    /// credential-bearing transport for a repo-chosen endpoint the user
118    /// has not explicitly trusted. The wrapped string is the actionable
119    /// message produced by [`crate::config::endpoint_credential_trust`].
120    #[error("{0}")]
121    UntrustedRemote(String),
122    /// A CAS ref write was rejected because the remote moved under us
123    /// (non-fast-forward). Callers map this to an actionable
124    /// fetch-then-retry / `--force-with-lease` hint.
125    #[error(
126        "updates were rejected for branch '{branch}' (non-fast-forward); fetch and merge first, or re-run with --force-with-lease / --force"
127    )]
128    NonFastForwardPush { branch: String },
129    /// The branch's packmap pointer could not be durably advanced under
130    /// sustained concurrent pushes. The branch ref was NOT moved, so the
131    /// remote stays consistent; the push is safe to retry.
132    #[error(
133        "could not establish the pack map for branch '{branch}' under concurrent pushes; retry"
134    )]
135    PackmapContended { branch: String },
136    /// A packmap chain was malformed — it exceeded the depth cap, contained
137    /// a cycle, or a node could not be downloaded/decoded. Indicates a
138    /// corrupt or hostile remote.
139    #[error("pack map chain for branch '{branch}' is malformed (too deep, cyclic, or unreadable)")]
140    PackChainInvalid { branch: String },
141    /// The remote advertised a branch ref but no packmap
142    /// (`refs/mkit/packmap/<branch>`) for it. mkit speaks a single,
143    /// packmap-only transfer dialect: the push path ALWAYS advertises a
144    /// packmap before moving the branch ref, so a branch with a tip but no
145    /// packmap is a corrupt/incomplete remote, not a format we degrade to.
146    /// We refuse to fetch rather than silently materialise a partial ref.
147    #[error("remote advertised branch '{0}' but no pack map to reconstruct it")]
148    PackmapMissing(String),
149    /// The packmap advertised a pack the remote does not hold. The branch's
150    /// closure cannot be reconstructed, so the fetch is aborted rather than
151    /// publishing a ref to an incomplete history.
152    #[error("remote advertised pack {pack} for branch '{branch}' but does not hold it")]
153    AdvertisedPackMissing { branch: String, pack: String },
154    /// After unpacking the branch's whole packmap chain, an object
155    /// reachable from the fetched tip is still absent from the local store —
156    /// the chain did not deliver the full closure. This is a pure integrity
157    /// assertion (no recovery download is attempted): the remote's packmap
158    /// is incomplete, so fetch aborts before publishing the ref.
159    #[error("remote is missing object {0} needed to reconstruct the ref")]
160    RemoteMissingObject(String),
161    /// The fetched tip's object closure exceeds the
162    /// [`mkit_core::ops::graph::MAX_REACHABLE`] verification cap, so
163    /// completeness could not be confirmed. On the applied-pack skip path
164    /// (#409) this closure walk is the sole guarantee the local store is
165    /// whole; a truncated walk could silently pass over missing objects, so
166    /// we fail closed rather than publish a ref we can't fully verify. This is
167    /// NOT a self-heal trigger.
168    #[error(
169        "fetched history is too large to verify (closure exceeds the {0}-object cap); refusing to publish an unverified ref"
170    )]
171    ClosureTooLarge(usize),
172    /// A commit/remix/tag newly introduced by this fetch failed Ed25519
173    /// signature verification via [`mkit_core::sign::verify_commit`] /
174    /// `verify_remix` / `verify_tag` — the exact check `mkit verify <rev>`
175    /// runs manually (issue #692). A hostile remote (THREAT-MODEL §3.1) can
176    /// otherwise push an unsigned or forged history that `clone`/`pull`/
177    /// `fetch` would silently materialise. Deliberately distinct from
178    /// [`RemoteMissingObject`](Self::RemoteMissingObject) so it does NOT
179    /// feed the applied-pack self-heal retry (#409): an invalid signature
180    /// is not evidence of local staleness, and clearing the applied-packs
181    /// record would not make a hostile remote's history valid. Fails
182    /// closed by default; opt out with `--no-verify-signatures` or the
183    /// user-scoped `pull.require_signed = false` config (never settable
184    /// from repo-scoped config — see [`crate::config::REPO_FORBIDDEN_KEYS`]).
185    #[error("object {hash} failed signature verification: {reason}")]
186    UnsignedOrInvalidObject { hash: String, reason: String },
187    /// A remote name passed to the applied-packs record (#409) is not a legal
188    /// ref name (per [`mkit_core::refs::validate_ref_name`]). A remote name
189    /// *is* a ref name, so this should never occur for a config-registered
190    /// remote; it is defence-in-depth against a malformed name being used as
191    /// a raw path component under `.mkit/applied-packs/`.
192    #[error("invalid remote name for applied-packs record: '{0}'")]
193    InvalidRemoteName(String),
194    /// One advance would need more data packs than the server permits per
195    /// advance. `commit` names the single commit or merge that cannot be split
196    /// further when the refusal came from the pre-flight dry seal, before
197    /// anything was published; without it the seal-time backstop fired.
198    #[error(
199        "push needs {packs} data packs but this server permits {limit} per advance{}{}; ask the operator to raise max_pack_bytes",
200        .commit.as_ref().map_or_else(String::new, |c| format!(" ({c} cannot be split further)")),
201        if *.holds_remote_head {
202            " (it is the first commit that contains the remote head: rebase onto the remote head, or merge it in a smaller or earlier commit)"
203        } else {
204            ""
205        }
206    )]
207    PushTooLarge {
208        packs: usize,
209        limit: usize,
210        commit: Option<String>,
211        /// The unsplittable commit is the first that contains the remote head,
212        /// so the split could not cut before it.
213        holds_remote_head: bool,
214    },
215    /// A retry of one advance (a rejected ticket, a lost delta base) re-planned
216    /// the push and the new plan no longer fits one advance. Found by a dry
217    /// seal, so nothing was uploaded for the retry.
218    #[error(
219        "{reason}, and the retry needs {packs} data packs but this server permits {limit} per advance; nothing was uploaded for the retry"
220    )]
221    RetryTooLarge {
222        reason: RetryReason,
223        packs: usize,
224        limit: usize,
225    },
226    /// A split push would exceed a bound (steps or chain length). Nothing was
227    /// published.
228    #[error("{0}")]
229    PushSplitLimit(String),
230    /// The stored grants do not authorize every advance of a split push.
231    /// Nothing was published.
232    #[error("{0}")]
233    PushNotAuthorized(String),
234    /// A split push failed after publishing some advances. `head` is the last
235    /// advance this push published, a first-parent ancestor of the local tip
236    /// (another pusher may have moved the branch since); re-running the push
237    /// resumes from the remote branch unless a concurrent push caused the
238    /// failure.
239    #[error(
240        "{cause}; {published} of {total} advances were published and the last published advance was {head} on branch '{branch}'{}",
241        resume_hint(.cause)
242    )]
243    SplitInterrupted {
244        branch: String,
245        head: String,
246        published: usize,
247        total: usize,
248        cause: Box<DispatchError>,
249    },
250    #[error("upload ticket was rejected after the push was retried")]
251    TicketRejected,
252    #[error("delta base is unavailable after a restart")]
253    DeltaBaseUnavailable,
254    #[error("packlist still names a pack absent from this repository after retry")]
255    PacklistNotInRepository,
256}
257
258/// The closing advice of a split push's failure: resume, unless the failure
259/// was a concurrent push (the branch then moved, and a re-run would only
260/// compare against whatever the other pusher left).
261fn resume_hint(cause: &DispatchError) -> &'static str {
262    if matches!(cause, DispatchError::NonFastForwardPush { .. }) {
263        "; the branch was moved by another push, so fetch and merge before pushing again"
264    } else {
265        "; re-run the push to resume"
266    }
267}
268
269/// Why an advance was re-planned.
270#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
271pub enum RetryReason {
272    /// The server no longer holds a delta base the plan relied on.
273    #[error(
274        "the server no longer holds a delta base this push relied on, so it must be re-sent as a full closure (ask the operator whether the repository was restored from an older copy)"
275    )]
276    LostDeltaBase,
277    /// The upload ticket or packlist was rejected, so the push was planned
278    /// again against the remote's current head, which may have moved.
279    #[error(
280        "the upload was rejected and the push was re-planned against the remote's current head, which may have moved (fetch and merge, then push again)"
281    )]
282    Replanned,
283}
284
285/// The advances a split push published before it failed.
286#[derive(Debug, Clone, PartialEq, Eq)]
287pub struct PublishedPrefix {
288    pub branch: String,
289    pub published: usize,
290    pub total: usize,
291    pub head: String,
292    /// Whether re-running the push is the advice (not after a concurrent push).
293    pub resumable: bool,
294}
295
296impl PublishedPrefix {
297    /// The note printed with any failure that left a prefix published.
298    #[must_use]
299    pub fn note(&self) -> String {
300        format!(
301            "{} of {} advances were published; the last published advance on branch '{}' was {}{}",
302            self.published,
303            self.total,
304            self.branch,
305            self.head,
306            if self.resumable {
307                "; re-run the push to resume"
308            } else {
309                "; the branch was moved by another push, so fetch and merge before pushing again"
310            }
311        )
312    }
313}
314
315impl DispatchError {
316    /// The error a split push failed with, without its published-prefix
317    /// wrapper, and that prefix (`None` when no advance of this branch was published).
318    #[must_use]
319    pub fn into_published_prefix(self) -> (Self, Option<PublishedPrefix>) {
320        match self {
321            Self::SplitInterrupted {
322                branch,
323                head,
324                published,
325                total,
326                cause,
327            } => {
328                let resumable = !matches!(*cause, Self::NonFastForwardPush { .. });
329                (
330                    *cause,
331                    Some(PublishedPrefix {
332                        branch,
333                        published,
334                        total,
335                        head,
336                        resumable,
337                    }),
338                )
339            }
340            other => (other, None),
341        }
342    }
343}
344
345/// Interpret a missing result as a repository failure only for repository-level
346/// operations. Pack downloads retain their content-specific missing errors.
347fn repository_operation_error(tx: &dyn Transport, error: TransportError) -> DispatchError {
348    if let TransportError::RemoteError(message) = &error
349        && message.starts_with("upload interrupted; ")
350    {
351        return DispatchError::UploadInterrupted(message.clone());
352    }
353    if matches!(&error, TransportError::RemoteError(message) if message == PENDING_INTERRUPTED_MESSAGE)
354    {
355        return DispatchError::Interrupted;
356    }
357    if matches!(&error, TransportError::PackNotFound)
358        && let Some(address) = tx.repository_address()
359    {
360        return DispatchError::RepositoryNotFound {
361            identity: address.repository.to_owned(),
362            origin: address.origin.to_owned(),
363        };
364    }
365    error.into()
366}
367
368/// Open a transport for `endpoint` only after the per-endpoint
369/// credential-trust gate (#97) approves it.
370///
371/// This is the single choke point through which push / fetch / pull
372/// (and named-remote callers in #175) MUST build a transport: it runs
373/// [`crate::config::endpoint_credential_trust`] — keyed on the resolved
374/// ENDPOINT and its `repo_chosen` provenance — *before* constructing
375/// the transport, so a credential-bearing HTTP/S3 transport is never
376/// instantiated for a repo-chosen endpoint the user hasn't trusted.
377///
378/// `repo_chosen` is `true` when the endpoint came from repo-scoped
379/// config (the flat `remote_endpoint` or a `remote.<name>.url`),
380/// `false` when it came from the user / an explicit CLI argument. Trust
381/// is per ENDPOINT, never per remote name.
382///
383/// `layout` is needed only to resolve a repo-key-file envelope signer
384/// when `cfg.merged.transport_auth == "envelope"` (see
385/// `envelope_signer_from_config`) — every caller already has it at
386/// hand (it discovered the repo before building `cfg`).
387pub fn open_trusted(
388    endpoint: &str,
389    remote_name: &str,
390    repo_chosen: bool,
391    cfg: &crate::config::LayeredConfig,
392    layout: &RepoLayout,
393) -> Result<Arc<dyn Transport>, DispatchError> {
394    crate::config::endpoint_credential_trust(cfg, endpoint, repo_chosen)
395        .map_err(DispatchError::UntrustedRemote)?;
396    open_with_config_for_remote(endpoint, &cfg.merged, layout, Some(remote_name))
397}
398
399/// A remote opened for `mkit push`: the transport, and (for a signing Connect
400/// client) the check that its stored grants authorize every advance of a split
401/// push.
402pub(crate) struct PushRemote {
403    pub tx: Arc<dyn Transport>,
404    pub authority: Option<Arc<dyn StepAuthority>>,
405}
406
407/// [`open_trusted`] for a push: the same trust gate and transport, plus the
408/// grant pre-check a split push needs (WP-1.17b, B5).
409///
410/// # Errors
411/// As [`open_trusted`].
412pub(crate) fn open_trusted_for_push(
413    endpoint: &str,
414    remote_name: &str,
415    repo_chosen: bool,
416    cfg: &crate::config::LayeredConfig,
417    layout: &RepoLayout,
418) -> Result<PushRemote, DispatchError> {
419    crate::config::endpoint_credential_trust(cfg, endpoint, repo_chosen)
420        .map_err(DispatchError::UntrustedRemote)?;
421    if !is_connect_url(endpoint) {
422        return Ok(PushRemote {
423            tx: open_with_config_for_remote(endpoint, &cfg.merged, layout, Some(remote_name))?,
424            authority: None,
425        });
426    }
427    let parts = open_connect_parts(endpoint, &cfg.merged, layout, Some(remote_name), true)?;
428    let tx = Arc::new(parts.transport);
429    let authority = parts
430        .signer_key
431        .zip(parts.grants)
432        .map(|(signer_key, grants)| {
433            Arc::new(grants::ConnectAuthority {
434                transport: tx.clone(),
435                grants,
436                signer_key,
437            }) as Arc<dyn StepAuthority>
438        });
439    Ok(PushRemote { tx, authority })
440}
441
442/// The single chokepoint that resolves SSH trust-pinning (issue #389) and
443/// `mkit+https://` envelope-signing config from `cfg` and opens a
444/// transport. Every config-bearing caller — [`open_trusted`] (push /
445/// fetch / pull) and `clone` — routes through here, so both are resolved
446/// and threaded in exactly ONE place. A new remote command physically
447/// cannot forget them as long as it opens through config; the only
448/// un-pinned path is the config-less [`open`], which production never
449/// uses for `ssh` or envelope auth.
450pub(crate) fn open_with_config(
451    url: &str,
452    cfg: &crate::config::Config,
453    layout: &RepoLayout,
454) -> Result<Arc<dyn Transport>, DispatchError> {
455    open_with_config_for_remote(url, cfg, layout, None)
456}
457
458fn open_with_config_for_remote(
459    url: &str,
460    cfg: &crate::config::Config,
461    layout: &RepoLayout,
462    remote_name: Option<&str>,
463) -> Result<Arc<dyn Transport>, DispatchError> {
464    if is_connect_url(url) {
465        return Ok(Arc::new(open_connect_with_config(
466            url,
467            cfg,
468            layout,
469            remote_name,
470            true,
471        )?));
472    }
473    open_with_ssh_options(url, &ssh_options_from_config(cfg), None)
474}
475
476fn is_connect_url(url: &str) -> bool {
477    url.starts_with("mkit+https://") || url.starts_with("mkit+http://")
478}
479
480/// [`open_trusted`] for a Connect endpoint, returning the concrete
481/// [`ConnectTransport`] so callers can reach the epoch and visibility RPCs
482/// (WP-2.14). `sign` chooses whether the ambient signing identity is used:
483/// `mkit epoch` and statement-mode `mkit visibility set` send unsigned RPCs
484/// and pass `false`; envelope-mode visibility passes `true` and so needs
485/// `transport_auth = envelope` and a trusted remote, like `mkit push`.
486///
487/// # Errors
488/// The credential-trust gate, a non-Connect endpoint, or connection setup.
489pub(crate) fn open_connect_trusted(
490    endpoint: &str,
491    remote_name: &str,
492    repo_chosen: bool,
493    cfg: &crate::config::LayeredConfig,
494    layout: &RepoLayout,
495    sign: bool,
496) -> Result<ConnectTransport, DispatchError> {
497    crate::config::endpoint_credential_trust(cfg, endpoint, repo_chosen)
498        .map_err(DispatchError::UntrustedRemote)?;
499    if !is_connect_url(endpoint) {
500        return Err(DispatchError::UnsupportedScheme(format!(
501            "`{endpoint}` is not an mkit+https:// or mkit+http:// remote; grant epochs and repository visibility are served over Connect"
502        )));
503    }
504    open_connect_with_config(endpoint, &cfg.merged, layout, Some(remote_name), sign)
505}
506
507/// Build the Connect transport for `url` from `cfg`: envelope signing (when
508/// `sign` and configured), the user grant store as the `GrantSource`, the
509/// upload receipt store, progress observers and the admission helper.
510pub(crate) fn open_connect_with_config(
511    url: &str,
512    cfg: &crate::config::Config,
513    layout: &RepoLayout,
514    remote_name: Option<&str>,
515    sign: bool,
516) -> Result<ConnectTransport, DispatchError> {
517    Ok(open_connect_parts(url, cfg, layout, remote_name, sign)?.transport)
518}
519
520/// A Connect transport with what a push needs to check its authority.
521struct ConnectParts {
522    transport: ConnectTransport,
523    /// The user grants installed on the transport (none without a signer).
524    grants: Option<Arc<grants::LocalGrants>>,
525    /// The signing key's public hex, when requests are signed.
526    signer_key: Option<String>,
527}
528
529fn open_connect_parts(
530    url: &str,
531    cfg: &crate::config::Config,
532    layout: &RepoLayout,
533    remote_name: Option<&str>,
534    sign: bool,
535) -> Result<ConnectParts, DispatchError> {
536    let envelope_signer = if sign {
537        if cfg.transport_auth_envelope() && cfg.trusted_remote_endpoint.trim() != url {
538            return Err(DispatchError::UntrustedRemote(format!(
539                "refusing request signing for untrusted destination `{url}`; run `mkit config trusted_remote_endpoint {url}` before using ambient signing identity"
540            )));
541        }
542        envelope_signer_from_config(cfg, layout)?
543    } else {
544        None
545    };
546    validate_connect_repository(url)?;
547    let signed = envelope_signer.is_some();
548    let signer_key = envelope_signer
549        .as_ref()
550        .map(|signer| signer.public_key_hex());
551    // Resolve only the selected config fallback: an authoritative environment
552    // path must not be defeated by an ignored fallback's tilde/HOME error.
553    let ca_file = if url.starts_with("mkit+https://")
554        && std::env::var_os(mkit_transport_connect::tls::CA_FILE_ENV).is_none()
555    {
556        cfg.ssl_ca_file_path(layout).map_err(|error| {
557            DispatchError::Transport(TransportError::TlsConfiguration(error.to_string()))
558        })?
559    } else {
560        None
561    };
562    let mut tx = ConnectTransport::connect_with_signer_and_ca_file(
563        url,
564        envelope_signer,
565        ca_file.as_deref(),
566    )?
567    .with_receipt_store(Arc::new(upload_receipts::FilePartReceiptStore::new(
568        layout.upload_parts_dir(),
569    )))
570    .with_pending_observer(|event| {
571        crate::progress::pending_event(event);
572        !crate::signal::is_shutdown()
573    })
574    .with_upload_observer(|event| {
575        crate::progress::upload_event(event);
576        !crate::signal::is_shutdown()
577    })
578    .with_admission_receipt_observer(|receipt| {
579        eprintln!(
580            "note: remote returned a {} receipt for {}",
581            receipt.header,
582            receipt
583                .procedure
584                .rsplit('/')
585                .next()
586                .unwrap_or(receipt.procedure)
587        );
588    });
589    // Grants ride only on signed requests (SPEC-WRITE-GRANTS §4.2), so the
590    // store is read only when there is a signer to present them with.
591    let grants = signed.then(|| Arc::new(grants::LocalGrants::from_stored(load_user_grants(cfg))));
592    if let Some(grants) = &grants {
593        tx = tx.with_grant_source(grants.clone());
594    }
595    if !cfg.admission_helper.is_empty() && cfg.trusted_remote_endpoint.trim() == url {
596        let responder = Arc::new(crate::admission_helper::ExecResponder {
597            path: cfg.admission_helper.clone().into(),
598        });
599        let mut policy = AdmissionPolicy::new(responder);
600        if let Some(name) = remote_name
601            && let Some(headers) = cfg.remote_admission_headers.get(name)
602        {
603            let bearer = std::env::var("MKIT_API_TOKEN").is_ok_and(|s| !s.is_empty());
604            for header in headers.split(',').map(str::trim).filter(|s| !s.is_empty()) {
605                if is_reserved(header, bearer) {
606                    eprintln!(
607                        "warning: ignoring reserved header `{header}` in remote.{name}.admission_headers (see SPEC-TRANSPORT-CONNECT §5.1)"
608                    );
609                } else {
610                    policy = policy.with_extra_allowed(header);
611                }
612            }
613        }
614        tx = tx.with_admission(policy);
615    }
616    Ok(ConnectParts {
617        transport: tx,
618        grants,
619        signer_key,
620    })
621}
622
623/// The verified grants in the user store, one warning per skipped file.
624fn load_user_grants(cfg: &crate::config::Config) -> Vec<crate::grants::store::StoredGrant> {
625    let rps = crate::grants::parse_relying_parties(&cfg.grant_webauthn_rp).unwrap_or_else(|e| {
626        eprintln!("warning: grant.webauthn_rp: {e}; treating no relying party as pinned");
627        Vec::new()
628    });
629    let store = match crate::grants::store::GrantStore::open_default() {
630        Ok(store) => store,
631        Err(e) => {
632            eprintln!("warning: {e}; using no stored grants");
633            return Vec::new();
634        }
635    };
636    let report = store.load(&rps);
637    for warning in &report.warnings {
638        eprintln!("warning: {warning}");
639    }
640    report.grants
641}
642
643/// The configured mkit signing key as an owner signer for `mkit grant`,
644/// `mkit epoch` and `mkit visibility`. Same key resolution as
645/// [`envelope_signer_from_config`], without its `transport_auth` gate: the
646/// owner asked for this signature explicitly.
647pub(crate) fn owner_ed25519_signer(
648    cfg: &crate::config::Config,
649    layout: &RepoLayout,
650) -> Result<Arc<dyn mkit_transport_connect::EnvelopeSigner>, String> {
651    match cfg.signer.as_str() {
652        "" | "legacy" => {
653            let key_path = crate::config::resolve_key_path(layout, &cfg.signing_key)
654                .map_err(|e| format!("signing_key: {e}"))?;
655            if !key_path.exists() {
656                return Err(format!(
657                    "no signing key at {} — run `mkit keygen` first",
658                    key_path.display()
659                ));
660            }
661            let kp = mkit_core::sign::load_key(&key_path).map_err(|e| format!("load key: {e}"))?;
662            Ok(Arc::new(envelope_signer::RepoKeyEnvelopeSigner::new(kp)))
663        }
664        "keystore" => Ok(Arc::new(envelope_signer::KeystoreEnvelopeSigner::open(
665            cfg,
666        )?)),
667        other => Err(format!(
668            "unknown signer `{other}` — expected `legacy` or `keystore`"
669        )),
670    }
671}
672
673/// Resolve an [`mkit_transport_connect::EnvelopeSigner`] from `cfg`, when
674/// `cfg.transport_auth_envelope()` is set — `Ok(None)` otherwise (the
675/// default: bearer-token-only, unchanged from #700/#701).
676///
677/// Reuses EXACTLY the same signer resolution as `mkit commit`'s
678/// [`crate::commands::commit::load_commit_signer`] (`cfg.signer` ==
679/// `""`/`"legacy"` -> the repo key file at `cfg.signing_key`; `"keystore"`
680/// -> `cfg.key.ed25519_ref_or_fallback()` via `mkit-keystore`) rather than
681/// inventing a parallel key path — the write envelope authenticates with
682/// the SAME Ed25519 identity that already signs the user's commits.
683///
684/// Both signer kinds sign the raw envelope digest directly (no
685/// SPEC-SIGNING commit/remix/tag domain prefix): the legacy path delegates
686/// to the EXISTING `mkit_attest::RepoKeySigner` (its `sign` already signs
687/// the given bytes directly — "the PAE's own `\"DSSEv1 \"` prefix is the
688/// domain separator" per its own doc comment — so no new raw-Ed25519 call
689/// site is needed here), the keystore path via `KeySigner::sign`, whose
690/// own contract already documents "Ed25519 signers return the 64-byte
691/// RFC 8032 signature over `msg`" — i.e. no domain digest applied, exactly
692/// what the envelope needs. See `envelope_signer.rs` for both adapters.
693pub(crate) fn envelope_signer_from_config(
694    cfg: &crate::config::Config,
695    layout: &RepoLayout,
696) -> Result<Option<Arc<dyn mkit_transport_connect::EnvelopeSigner>>, DispatchError> {
697    if !cfg.transport_auth_envelope() {
698        return Ok(None);
699    }
700    let remote_error = |msg: String| DispatchError::Transport(TransportError::RemoteError(msg));
701    match cfg.signer.as_str() {
702        "" | "legacy" => {
703            let key_path =
704                crate::config::resolve_key_path(layout, &cfg.signing_key).map_err(|e| {
705                    remote_error(format!("transport_auth = envelope: signing_key: {e}"))
706                })?;
707            if !key_path.exists() {
708                return Err(remote_error(format!(
709                    "transport_auth = envelope requires a signing key at {} — run `mkit keygen` first",
710                    key_path.display()
711                )));
712            }
713            let kp = mkit_core::sign::load_key(&key_path)
714                .map_err(|e| remote_error(format!("transport_auth = envelope: load key: {e}")))?;
715            Ok(Some(
716                Arc::new(envelope_signer::RepoKeyEnvelopeSigner::new(kp))
717                    as Arc<dyn mkit_transport_connect::EnvelopeSigner>,
718            ))
719        }
720        "keystore" => {
721            let signer = envelope_signer::KeystoreEnvelopeSigner::open(cfg)
722                .map_err(|e| remote_error(format!("transport_auth = envelope: {e}")))?;
723            Ok(Some(
724                Arc::new(signer) as Arc<dyn mkit_transport_connect::EnvelopeSigner>
725            ))
726        }
727        other => Err(remote_error(format!(
728            "transport_auth = envelope: unknown signer `{other}` — expected `legacy` or `keystore`"
729        ))),
730    }
731}
732
733/// Map the three `ssh.*` trust-pinning keys from a merged [`Config`] into
734/// the [`SshOptions`] carried to the spawned `ssh(1)` child. An empty
735/// string means "unset" — `build_ssh_command` emits no flag for it, so
736/// the user's `ssh(1)` defaults are inherited. The producer half of
737/// issue #389 (the consumer half, `build_ssh_command`, wires the fields
738/// into argv). Sole caller is [`open_with_config`].
739fn ssh_options_from_config(cfg: &crate::config::Config) -> SshOptions {
740    SshOptions {
741        strict_host_key_checking: cfg.ssh_strict_host_key_checking.clone(),
742        user_known_hosts_file: cfg.ssh_user_known_hosts_file.clone(),
743        identity_file: cfg.ssh_identity_file.clone(),
744    }
745}
746
747/// Open a transport for the given URL with **no** SSH trust-pinning.
748/// Returns a type-erased `Arc` so callers can treat all schemes
749/// uniformly.
750///
751/// Low-level scheme dispatch only — it neither enforces the credential
752/// gate nor threads `ssh.*` config. Any caller that has a [`Config`]
753/// must use `open_with_config` (directly, or via `open_trusted`) so
754/// the trust-pinning keys reach the spawned `ssh(1)`; `open` stays
755/// public only for file/memory integration tests that have no ambient
756/// config to resolve.
757///
758/// [`Config`]: crate::config::Config
759pub fn open(url: &str) -> Result<Arc<dyn Transport>, DispatchError> {
760    open_with_ssh_options(url, &SshOptions::default(), None)
761}
762
763/// Scheme dispatch with explicit SSH options and an optional `mkit+https://`
764/// / `mkit+http://` envelope signer. Identical to [`open`] for every
765/// non-SSH, non-Connect scheme; the `mkit+ssh://` branch threads
766/// `ssh_options` (issue #389) into the spawned `ssh(1)` child via
767/// [`SshTransport::connect_with_options`], and the `mkit+https://`/
768/// `mkit+http://` branch threads `envelope_signer` (issue #699 follow-up)
769/// into [`ConnectTransport::connect_with_signer`]. Reached via [`open`]
770/// (no config — both `None`/default) and [`open_with_config`] for non-Connect
771/// schemes; `open_with_config` handles Connect URLs itself, so the Connect
772/// branch here is reached only from [`open`].
773fn open_with_ssh_options(
774    url: &str,
775    ssh_options: &SshOptions,
776    envelope_signer: Option<Arc<dyn mkit_transport_connect::EnvelopeSigner>>,
777) -> Result<Arc<dyn Transport>, DispatchError> {
778    if url.starts_with("git+") {
779        return Err(DispatchError::UnsupportedScheme(format!(
780            "'{url}' is a git-bridge remote — native push/pull/fetch/clone do not \
781             speak git transports; use `mkit git export` / `mkit git import` / \
782             `mkit git pull` (feature git-bridge)"
783        )));
784    }
785    if let Some(rest) = url.strip_prefix("mkit+file://") {
786        // mkit+file:///abs/path -> /abs/path
787        let path = Path::new(rest);
788        return Ok(Arc::new(FileTransport::new(path)));
789    }
790    if url.starts_with("mkit+memory://") {
791        // Memory transport is in-process; the URL-based path is not
792        // useful on its own but we accept it so `mkit remote add`
793        // round-trips cleanly.
794        return Err(DispatchError::UnsupportedScheme(
795            "mkit+memory:// must be driven via in-process harness (see tests)".to_string(),
796        ));
797    }
798    if url.starts_with("mkit+https://") || url.starts_with("mkit+http://") {
799        validate_connect_repository(url)?;
800        // ConnectTransport::connect_with_signer strips the `mkit+` prefix
801        // itself and reads MKIT_API_TOKEN from the environment (mkit#701 —
802        // the native mkit.transport.v1 ConnectRPC client, replacing the
803        // retired mkit-transport-http JSON dialect as of
804        // SPEC-TRANSPORT-CONNECT verb parity). Only the config-less `open`
805        // reaches this branch (`open_with_config` returns early for Connect
806        // URLs), so `envelope_signer` is always `None` here; bearer token and
807        // envelope signing are independent, additive auth modes.
808        let tx = ConnectTransport::connect_with_signer(url, envelope_signer)?;
809        return Ok(Arc::new(tx));
810    }
811    if url.starts_with("mkit+s3://") {
812        // S3Transport::connect reads MKIT_R2_ACCESS_KEY_ID /
813        // MKIT_R2_SECRET_ACCESS_KEY from the environment. Missing
814        // credentials surface as AccessDenied on the first signed call,
815        // not at connect time.
816        let tx = S3Transport::connect(url)?;
817        return Ok(Arc::new(tx));
818    }
819    if url.starts_with("mkit+ssh://") {
820        // Parse the URL, then spawn `ssh(1)` with the caller-supplied
821        // trust-pinning options (issue #389). `connect_with_options`
822        // performs the `Hello` / `HelloResponse` handshake. Any failure
823        // here tears the child down before returning, so callers never
824        // see a half-initialised transport.
825        let target = parse_mkit_ssh_url(url).map_err(SshInitError::from)?;
826        let tx = SshTransport::connect_with_options(&target, ssh_options)?;
827        return Ok(Arc::new(tx));
828    }
829    #[cfg(feature = "enc-transport")]
830    if url.starts_with("mkit+enc://") {
831        return open_enc(url);
832    }
833    Err(DispatchError::MalformedUrl(url.to_string()))
834}
835
836fn validate_connect_repository(url: &str) -> Result<(), DispatchError> {
837    repository_identity_from_url(url).map_err(|reason| {
838        let path = url
839            .split_once("://")
840            .and_then(|(_, rest)| rest.split_once('/'))
841            .map_or("", |(_, path)| path)
842            .split(['?', '#'])
843            .next()
844            .unwrap_or("")
845            .trim_matches('/');
846        DispatchError::MalformedUrl(format!("repository identity `{path}` in {url}: {reason}"))
847    })?;
848    Ok(())
849}
850
851/// `mkit+enc://` dispatch (issue #156).
852///
853/// Parses the URL, derives an ephemeral dialer keypair (keystore
854/// integration is SPEC-TRANSPORT-ENC §6 item 5, still deferred), and
855/// runs the encrypted-stream handshake against the URL-advertised
856/// server public key.
857///
858/// Client identity (issue #178): an allowlisting server pins the
859/// dialer's static ed25519 key. To survive across restarts the client
860/// can supply a STABLE raw-32 key file via the `MKIT_ENC_CLIENT_KEY`
861/// environment variable (a user-scoped / CLI-supplied path — never
862/// repo-local `.mkit/config`, which `open_enc` has no access to anyway).
863/// When the variable is unset we fall back to a fresh ephemeral key per
864/// process, which still works against an allow-any enc listener.
865#[cfg(feature = "enc-transport")]
866const ENC_CLIENT_KEY_ENV: &str = "MKIT_ENC_CLIENT_KEY";
867
868#[cfg(feature = "enc-transport")]
869fn open_enc(url: &str) -> Result<Arc<dyn Transport>, DispatchError> {
870    use mkit_transport_enc::url::parse_enc_url;
871
872    let target = parse_enc_url(url).map_err(DispatchError::Transport)?;
873    let sk = load_or_ephemeral_client_key()?;
874    let tx = mkit_transport_enc::connect_tcp(&target.host, target.port, &target.server_pubkey, sk)
875        .map_err(|e| DispatchError::Transport(TransportError::RemoteError(e.to_string())))?;
876    Ok(Arc::new(tx))
877}
878
879/// Resolve the dialer's static signing key.
880///
881/// If `MKIT_ENC_CLIENT_KEY` points at a raw 32-byte key file, load it
882/// (with the standard `load_raw_32` 0600/owner hardening) so the
883/// client's public key is stable — letting an allowlisting server pin
884/// it across restarts. Otherwise draw a fresh ephemeral key from the
885/// system RNG (≥256 bits) for back-compat with allow-any servers.
886#[cfg(feature = "enc-transport")]
887fn load_or_ephemeral_client_key()
888-> Result<commonware_cryptography::ed25519::PrivateKey, DispatchError> {
889    use commonware_codec::DecodeExt as _;
890    use commonware_cryptography::ed25519::PrivateKey;
891    use zeroize::Zeroizing;
892
893    let map_err = |e: String| DispatchError::Transport(TransportError::RemoteError(e));
894
895    if let Some(path) = std::env::var_os(ENC_CLIENT_KEY_ENV).filter(|s| !s.is_empty()) {
896        let seed = mkit_core::sign::load_raw_32(std::path::Path::new(&path))
897            .map_err(|e| map_err(format!("load {ENC_CLIENT_KEY_ENV}: {e}")))?;
898        return PrivateKey::decode(seed.as_ref())
899            .map_err(|e| map_err(format!("client key construction failed: {e}")));
900    }
901
902    // Ephemeral fallback. Draw 32 bytes from `getrandom`, wrapped in
903    // `Zeroizing` so the stack copy is scrubbed on drop; the resulting
904    // `PrivateKey` carries its own `Secret`-based zeroization.
905    let mut secret = Zeroizing::new([0u8; 32]);
906    getrandom::fill(secret.as_mut()).map_err(|e| map_err(e.to_string()))?;
907    PrivateKey::decode(secret.as_ref()).map_err(|e| map_err(e.to_string()))
908}
909
910/// Push every ref under `refs/heads/` to the remote. Returns the count of
911/// refs pushed. Each branch is published with [`push_branch`], which sends
912/// one delta-compressed pack of the objects the remote lacks, advertises it
913/// via the `refs/mkit/packmap/<branch>` ref, then moves the branch ref.
914pub fn push_all(cwd: &Path, tx: &dyn Transport) -> Result<usize, DispatchError> {
915    push_all_with(cwd, tx, None, false, None).map(|pushed| pushed.refs)
916}
917
918/// What a successful push published.
919#[derive(Debug, Clone, Copy, PartialEq, Eq)]
920pub struct Pushed {
921    /// Branches pushed.
922    pub refs: usize,
923    /// Advances published in total (one per branch unless a push was split).
924    pub steps: usize,
925}
926
927/// CAS-aware mirror push (`mkit push --all`). Pushes every local
928/// `refs/heads/*` to the remote, using the remote-tracking ref under
929/// `refs/remotes/<remote>/<branch>` as the CAS lease (Missing when no
930/// tracking ref exists, Match otherwise). `force` upgrades every write
931/// to an unconditional `Any`. On success each pushed branch's
932/// remote-tracking ref is advanced to the pushed tip.
933///
934/// `remote` is the remote NAME used for the local tracking-ref
935/// namespace; `None` means the legacy `default`. A branch too large for one
936/// advance is split (see [`push_branch_steps`]); its tracking ref follows each
937/// published advance, so an interrupted push resumes where it stopped.
938/// `authority` is checked before a split branch uploads anything.
939pub fn push_all_with(
940    cwd: &Path,
941    tx: &dyn Transport,
942    remote: Option<&str>,
943    force: bool,
944    authority: Option<&dyn StepAuthority>,
945) -> Result<Pushed, DispatchError> {
946    let layout = mkit_core::layout::discover(cwd)?;
947    let store = crate::commands::open_store_configured(&layout)?;
948    let refs_list = crate::commands::list_refs_parallel(&layout)?;
949    let remote = remote.unwrap_or(DEFAULT_REMOTE);
950    let mut n = 0;
951    let mut steps = 0;
952    let shallow = shallow_boundaries(&layout)?;
953    // Batch every pushed branch's remote-tracking-ref write (#645):
954    // publishing lands each ref as soon as its branch's `push_branch`
955    // succeeds (same visibility as before), but the directory fsync that
956    // makes those renames crash-durable is deferred to one pass over the
957    // distinct directories touched, below — instead of once per branch.
958    let mut tracking = refs::RemoteRefBatch::new(&layout, remote)?;
959    let result: Result<(), DispatchError> = (|| {
960        for r in refs_list {
961            if crate::signal::is_shutdown() {
962                return Err(DispatchError::Interrupted);
963            }
964            let Some(h) = r.hash else { continue };
965            let condition = if force {
966                refs::RefWriteCondition::Any
967            } else {
968                match refs::read_remote_ref(&layout, remote, &r.name)? {
969                    Some(tracked) => refs::RefWriteCondition::Match(tracked),
970                    None => refs::RefWriteCondition::Missing,
971                }
972            };
973            let control = PushControl {
974                authority,
975                shallow: shallow.clone(),
976                ..PushControl::default()
977            };
978            steps += push_branch_steps(
979                tx,
980                &store,
981                &r.name,
982                h,
983                condition,
984                rebaseline_depth(),
985                pack::MAX_TOTAL_PAYLOAD,
986                &control,
987                &mut |published| Ok(tracking.write(&r.name, &published)?),
988            )?;
989            n += 1;
990        }
991        Ok(())
992    })();
993    // Commit whatever tracking-ref writes succeeded regardless of how the
994    // loop above ended, so a mid-loop failure still durably publishes the
995    // prefix that already pushed successfully — matching the old
996    // per-branch loop, where each completed ref write was independently
997    // durable before the loop moved to the next branch.
998    tracking.commit()?;
999    result?;
1000    Ok(Pushed { refs: n, steps })
1001}
1002
1003fn shallow_boundaries(
1004    layout: &RepoLayout,
1005) -> Result<std::collections::HashSet<Hash>, DispatchError> {
1006    Ok(refs::load_shallow_boundaries(layout)?
1007        .unwrap_or_default()
1008        .into_iter()
1009        .collect())
1010}
1011
1012/// True iff advancing a ref from `old` to `new` is a fast-forward (i.e.
1013/// `old` is an ancestor of `new`). A missing `old` (brand-new ref) and an
1014/// unchanged ref both count as fast-forwards. Used by `push`/`fetch` to
1015/// pick the git-style summary symbol (`..` vs `...`/`(forced update)`).
1016pub fn is_fast_forward(cwd: &Path, old: Option<Hash>, new: Hash) -> Result<bool, DispatchError> {
1017    match old {
1018        None => Ok(true),
1019        Some(o) if o == new => Ok(true),
1020        Some(o) => {
1021            let layout = mkit_core::layout::discover(cwd)?;
1022            let store = crate::commands::open_store_configured(&layout)?;
1023            Ok(is_ancestor(&store, o, new)?)
1024        }
1025    }
1026}
1027
1028/// CAS lease policy for a default (current-branch → upstream) push.
1029#[derive(Debug, Clone, Copy)]
1030pub enum PushLease {
1031    /// Force — unconditional `Any`.
1032    Force,
1033    /// `--force-with-lease` — require the remote tip to equal the local
1034    /// remote-tracking ref (Match), or Missing when there is none.
1035    /// Identical mechanism to the default safe push; semantically it is
1036    /// the explicit, opt-in form that overwrites a fast-forward-failing
1037    /// branch *only* if the remote hasn't moved past what we last saw.
1038    WithLease,
1039    /// Default safe push: Match the local remote-tracking ref, or
1040    /// Missing when absent (first push of this branch).
1041    FastForward,
1042}
1043
1044/// Resolve the CAS condition for a single-branch push from the local
1045/// remote-tracking ref `refs/remotes/<remote>/<branch>` and the lease
1046/// policy.
1047pub fn lease_condition(
1048    cwd: &Path,
1049    remote: &str,
1050    branch: &str,
1051    lease: PushLease,
1052) -> Result<refs::RefWriteCondition, DispatchError> {
1053    if matches!(lease, PushLease::Force) {
1054        return Ok(refs::RefWriteCondition::Any);
1055    }
1056    let layout = mkit_core::layout::discover(cwd)?;
1057    Ok(match refs::read_remote_ref(&layout, remote, branch)? {
1058        Some(tracked) => refs::RefWriteCondition::Match(tracked),
1059        None => refs::RefWriteCondition::Missing,
1060    })
1061}
1062
1063/// Push the current branch to its upstream and, on success, advance the
1064/// local remote-tracking ref `refs/remotes/<remote>/<branch>` to the
1065/// pushed tip.
1066///
1067/// `remote` is the upstream remote NAME (for the tracking-ref
1068/// namespace); `branch` is the local branch name; `remote_branch` is the
1069/// branch name on the remote (`refs/heads/<remote_branch>`).
1070///
1071/// Returns the pushed tip and the number of advances it took: more than one
1072/// when the push was split along first-parent history (see
1073/// [`push_branch_steps`]). The tracking ref follows every published advance.
1074pub fn push_branch_tracked(
1075    cwd: &Path,
1076    tx: &dyn Transport,
1077    remote: &str,
1078    branch: &str,
1079    remote_branch: &str,
1080    lease: PushLease,
1081    authority: Option<&dyn StepAuthority>,
1082) -> Result<(Hash, usize), DispatchError> {
1083    let layout = mkit_core::layout::discover(cwd)?;
1084    let store = crate::commands::open_store_configured(&layout)?;
1085    let tip = refs::read_ref(&layout, branch)?
1086        .ok_or_else(|| DispatchError::RemoteBranchMissing(branch.to_owned()))?;
1087    // Default safe push requires a TRUE fast-forward: the new tip must
1088    // descend from the last-seen remote-tracking ref. The CAS `Match`
1089    // lease alone only proves the remote hasn't moved since we last
1090    // fetched — on its own it would still let a divergent local tip
1091    // (e.g. after a local `reset` to an unrelated commit) overwrite the
1092    // remote, which Git rejects as non-fast-forward. `--force-with-lease`
1093    // (`WithLease`) intentionally skips this check (overwrite as long as
1094    // the remote matches what we last saw); `Force` skips everything.
1095    if matches!(lease, PushLease::FastForward)
1096        && let Some(tracked) = refs::read_remote_ref(&layout, remote, remote_branch)?
1097        && !is_ancestor(&store, tracked, tip)?
1098    {
1099        return Err(DispatchError::NonFastForwardPush {
1100            branch: remote_branch.to_owned(),
1101        });
1102    }
1103    let condition = lease_condition(cwd, remote, remote_branch, lease)?;
1104    let control = PushControl {
1105        authority,
1106        shallow: shallow_boundaries(&layout)?,
1107        ..PushControl::default()
1108    };
1109    let steps = push_branch_steps(
1110        tx,
1111        &store,
1112        remote_branch,
1113        tip,
1114        condition,
1115        rebaseline_depth(),
1116        pack::MAX_TOTAL_PAYLOAD,
1117        &control,
1118        &mut |published| {
1119            Ok(refs::write_remote_ref(
1120                &layout,
1121                remote,
1122                remote_branch,
1123                &published,
1124            )?)
1125        },
1126    )?;
1127    Ok((tip, steps))
1128}
1129
1130/// Push one branch: upload one or more delta-compressed packs carrying
1131/// every object reachable from `tip` that the remote lacks — split
1132/// across multiple packs when the plan's payload exceeds a single
1133/// pack's cap (issue #831) — durably advertise them as one node on the
1134/// `refs/mkit/packmap/<branch>` metadata ref, then CAS-write
1135/// `refs/heads/<branch>` under `condition`.
1136///
1137/// Objects already present at the remote's current tip are never re-sent
1138/// (identical-object dedup), and changed `FastCDC` chunks are delta-encoded
1139/// against the prior version the remote already holds when that saves
1140/// bytes (see [`mkit_core::transfer::plan_pack`]). The pack is keyed by its
1141/// own BLAKE3 digest (SPEC-PACKFILE §7) — required because the digest-
1142/// checking storage server rejects a delta stored under the reconstructed
1143/// object's hash.
1144///
1145/// The packmap is advanced *and confirmed* before the branch ref moves: if
1146/// the packmap can't be durably established the push aborts without
1147/// touching the head, so the head never points past a packmap that fails to
1148/// reconstruct it (even under concurrent pushers to the same branch).
1149///
1150/// Plans the pack FIRST (diffing against the remote's current tip). A no-op
1151/// push (empty plan — the remote already holds this closure) takes the cheap
1152/// head-only path and walks NO packmap chain (mkit #521 perf). Only when the
1153/// plan is non-empty does it resolve the branch's current packmap chain depth
1154/// (walking it exactly once, see `packmap::probe_chain`) and, if the chain
1155/// would grow past the re-baseline threshold (#406, see
1156/// `packmap::rebaseline_depth`) AND the transport's `advance_refs` is
1157/// transactional ([`Transport::supports_atomic_advance`], mkit #521) AND the
1158/// head write is CAS-conditioned (a force push's `Any` head condition takes
1159/// the safe append path — an `Any` condition makes even an atomic transport
1160/// fall back to the ordered two-PUT `advance_refs`, so a reset there is not
1161/// safe), re-plans as a full closure (diffs against no remote tip) and
1162/// carries that decision down to `advance_packmap` as
1163/// `ChainAction::ResetSelfContained` so it resets the chain to a single
1164/// fresh node instead of appending to it — bounding clone cost, which
1165/// otherwise grows with chain length.
1166///
1167/// On a transport WITHOUT transactional `advance_refs` (the default used by
1168/// file/S3/SSH/memory), crossing the threshold never triggers a reset: the
1169/// default `advance_refs` commits the packmap write before the head CAS, and
1170/// a reset (unlike an append) is not a superset of the prior chain, so a
1171/// lost head-CAS race after a committed reset would strand the (unmoved)
1172/// head pointing at a commit the packmap can no longer reconstruct. Such a
1173/// transport keeps appending — `ChainAction::Append` — past the
1174/// threshold; `packmap::MAX_PACK_CHAIN_DEPTH` (the pure runaway/cycle guard)
1175/// remains the only bound on chain growth there, unchanged by this gate.
1176///
1177/// A chain read that fails with [`DispatchError::PackChainInvalid`], or a
1178/// missing prior packmap (first push), is left alone here — depth is only
1179/// defined for a resolvable chain, and a broken chain already has its own
1180/// reset path in `advance_packmap` (the broken-chain escape hatch, gated on
1181/// `self_contained` alone, independent of this transactional-advance gate —
1182/// see `ChainAction::Append`'s doc comment).
1183///
1184/// The already-resolved chain from this probe (when not discarded by a
1185/// re-baseline decision) is threaded into `advance_packmap` so its first
1186/// CAS attempt does not have to walk the chain a second time (#521 perf
1187/// fix).
1188///
1189/// On a CAS failure ([`TransportError::RefConflict`]) this reads the head
1190/// back first (SPEC-TRANSPORT §7: a retried write can report a conflict for
1191/// its own landed attempt); a head already at `tip` is success, anything
1192/// else returns [`DispatchError::NonFastForwardPush`] so callers can render
1193/// an actionable fetch-then-retry hint. Does NOT touch local
1194/// remote-tracking refs — the caller decides when to advance them.
1195pub fn push_branch(
1196    tx: &dyn Transport,
1197    store: &ObjectStore,
1198    branch: &str,
1199    tip: Hash,
1200    condition: refs::RefWriteCondition,
1201) -> Result<(), DispatchError> {
1202    push_branch_with_depth(tx, store, branch, tip, condition, rebaseline_depth())
1203}
1204
1205/// [`push_branch`] with an explicit re-baseline threshold in place of the
1206/// configured one (`packmap::rebaseline_depth`: the
1207/// `MKIT_PACK_REBASELINE_DEPTH` env var, default 64). `0` disables
1208/// re-baselining. Semantics are otherwise identical — see [`push_branch`].
1209///
1210/// This is the depth-policy seam (#547): integration tests inject a small
1211/// threshold (e.g. 3) to exercise the re-baseline path in-process with a
1212/// handful of pushes, where reaching the default threshold would take ~64
1213/// real pushes and the env-var override cannot be set on the test's own
1214/// process (`std::env::set_var` is banned by `clippy::disallowed_methods` —
1215/// it races other threads on POSIX).
1216pub fn push_branch_with_depth(
1217    tx: &dyn Transport,
1218    store: &ObjectStore,
1219    branch: &str,
1220    tip: Hash,
1221    condition: refs::RefWriteCondition,
1222    rebaseline_threshold: usize,
1223) -> Result<(), DispatchError> {
1224    push_branch_with_limits(
1225        tx,
1226        store,
1227        branch,
1228        tip,
1229        condition,
1230        rebaseline_threshold,
1231        pack::MAX_TOTAL_PAYLOAD,
1232    )
1233}
1234
1235/// [`push_branch_with_depth`] with an explicit per-pack payload cap in
1236/// place of the format's hardcoded [`pack::MAX_TOTAL_PAYLOAD`] (issue
1237/// #831). Semantics are otherwise identical — see [`push_branch`].
1238///
1239/// This is the payload-cap test seam, mirroring the #547
1240/// `rebaseline_threshold` pattern: integration tests inject a tiny cap
1241/// (a few KiB) to exercise multi-pack splitting in-process without
1242/// moving a multi-GiB plan through it.
1243pub fn push_branch_with_limits(
1244    tx: &dyn Transport,
1245    store: &ObjectStore,
1246    branch: &str,
1247    tip: Hash,
1248    condition: refs::RefWriteCondition,
1249    rebaseline_threshold: usize,
1250    pack_payload_cap: u64,
1251) -> Result<(), DispatchError> {
1252    push_branch_steps(
1253        tx,
1254        store,
1255        branch,
1256        tip,
1257        condition,
1258        rebaseline_threshold,
1259        pack_payload_cap,
1260        &PushControl::default(),
1261        &mut |_| Ok(()),
1262    )
1263    .map(drop)
1264}
1265
1266/// [`push_branch_with_limits`] that splits a push too large for one advance
1267/// along the branch's first-parent history (WP-1.17b, R-164), and reports each
1268/// published advance to `on_step` with the commit the remote branch now
1269/// points at. Returns the number of advances.
1270///
1271/// The split applies only when the transport ticketed-uploads
1272/// ([`UploadLimits::tickets_per_advance`]) and the conservative estimate of the
1273/// push's packs exceeds the data-pack budget of one advance. Otherwise this is
1274/// exactly one [`push_branch_with_limits`] call. When it applies, the first
1275/// advance uses `condition`, and every later one is a `Match` on the previous
1276/// advance's commit, whatever `condition` was, so a long split never
1277/// overwrites a concurrent pusher. Each advance uploads its own packs just
1278/// before it lands; nothing is pre-uploaded.
1279///
1280/// A failure after the first advance is returned as
1281/// [`DispatchError::SplitInterrupted`], naming the published prefix; running
1282/// the push again resumes from it.
1283///
1284/// # Errors
1285/// As [`push_branch`], plus the split's own refusals (see
1286/// [`plan_push_steps`]), and whatever `on_step` returns.
1287#[allow(clippy::too_many_arguments)]
1288pub fn push_branch_steps(
1289    tx: &dyn Transport,
1290    store: &ObjectStore,
1291    branch: &str,
1292    tip: Hash,
1293    condition: refs::RefWriteCondition,
1294    rebaseline_threshold: usize,
1295    pack_payload_cap: u64,
1296    control: &PushControl<'_>,
1297    on_step: &mut dyn FnMut(Hash) -> Result<(), DispatchError>,
1298) -> Result<usize, DispatchError> {
1299    let limits = tx.upload_limits();
1300    let mut single = |plan| {
1301        push_step_recovering(
1302            tx,
1303            store,
1304            branch,
1305            tip,
1306            condition,
1307            rebaseline_threshold,
1308            pack_payload_cap,
1309            plan,
1310        )?;
1311        on_step(tip)?;
1312        Ok(1)
1313    };
1314    if limits.tickets_per_advance.is_none() {
1315        return single(None);
1316    }
1317    refs::check_pushable_branch(branch)?;
1318    let remote_tip = tx.read_ref(&format!("refs/heads/{branch}"))?;
1319    let plan = transfer::plan_pack_with(store, tip, remote_tip, encode_delta_candidates_batch)?;
1320    let cap = effective_payload_cap(pack_payload_cap, limits.max_pack_bytes)?;
1321    if estimate_pack_sizes(store, &plan, cap, limits.max_pack_bytes)?.len()
1322        <= split::data_pack_budget(limits)
1323    {
1324        return single(Some(plan));
1325    }
1326    // The estimate is uncompressed: a push that compresses into the budget
1327    // is not split (exact dry seal), and one against a remote tip this store
1328    // lacks is left to the head CAS rather than walking all of history.
1329    let unknown_remote = remote_tip.is_some_and(|remote| !store.contains(&remote));
1330    if unknown_remote {
1331        return single(None);
1332    }
1333    match build_and_upload_packs(PackSink::Count, store, plan, cap, limits) {
1334        Err(DispatchError::PushTooLarge { .. }) => {}
1335        Err(DispatchError::Interrupted) => return Err(DispatchError::Interrupted),
1336        _ => return single(None),
1337    }
1338    let steps = plan_push_steps(
1339        store,
1340        tip,
1341        remote_tip,
1342        limits,
1343        pack_payload_cap,
1344        control,
1345        condition,
1346        branch,
1347    )?;
1348    if steps.len() == 1 {
1349        return single(None);
1350    }
1351    let total = steps.len();
1352    let mut published: Option<Hash> = None;
1353    for (index, step_tip) in steps.into_iter().enumerate() {
1354        let step_condition = published.map_or(condition, refs::RefWriteCondition::Match);
1355        let mut landed = false;
1356        let outcome = (|| {
1357            if crate::signal::is_shutdown() {
1358                return Err(DispatchError::Interrupted);
1359            }
1360            crate::progress::report(crate::progress::Event::Step {
1361                index: index + 1,
1362                total,
1363            });
1364            push_step_recovering(
1365                tx,
1366                store,
1367                branch,
1368                step_tip,
1369                step_condition,
1370                rebaseline_threshold,
1371                pack_payload_cap,
1372                None,
1373            )?;
1374            landed = true;
1375            on_step(step_tip)?;
1376            crate::progress::step_committed(index + 1, total, &step_tip);
1377            Ok(())
1378        })();
1379        if landed {
1380            published = Some(step_tip);
1381        }
1382        if let Err(cause) = outcome {
1383            return Err(match published {
1384                Some(head) => DispatchError::SplitInterrupted {
1385                    branch: branch.to_owned(),
1386                    head: mkit_core::hash::to_hex(&head),
1387                    published: index + usize::from(landed),
1388                    total,
1389                    cause: Box::new(cause),
1390                },
1391                None => cause,
1392            });
1393        }
1394    }
1395    Ok(total)
1396}
1397
1398/// One advance to `tip` with the restart-once recovery: a rejected ticket or
1399/// packlist re-plans the same way, a missing delta base re-plans as the full
1400/// closure. `plan` is an already-computed plan against the remote's current
1401/// tip, used for the first attempt only.
1402#[allow(clippy::too_many_arguments)]
1403fn push_step_recovering(
1404    tx: &dyn Transport,
1405    store: &ObjectStore,
1406    branch: &str,
1407    tip: Hash,
1408    condition: refs::RefWriteCondition,
1409    rebaseline_threshold: usize,
1410    pack_payload_cap: u64,
1411    plan: Option<transfer::PackPlan>,
1412) -> Result<(), DispatchError> {
1413    let once = |retry, plan| {
1414        push_branch_once(
1415            tx,
1416            store,
1417            branch,
1418            tip,
1419            condition,
1420            rebaseline_threshold,
1421            pack_payload_cap,
1422            retry,
1423            plan,
1424        )
1425    };
1426    match once(None, plan) {
1427        Err(DispatchError::TicketRejected | DispatchError::PacklistNotInRepository) => {
1428            once(Some(RetryReason::Replanned), None)
1429        }
1430        Err(DispatchError::DeltaBaseUnavailable) => once(Some(RetryReason::LostDeltaBase), None),
1431        result => result,
1432    }
1433}
1434
1435#[allow(clippy::too_many_arguments)]
1436fn push_branch_once(
1437    tx: &dyn Transport,
1438    store: &ObjectStore,
1439    branch: &str,
1440    tip: Hash,
1441    condition: refs::RefWriteCondition,
1442    rebaseline_threshold: usize,
1443    pack_payload_cap: u64,
1444    retry: Option<RetryReason>,
1445    preplanned: Option<transfer::PackPlan>,
1446) -> Result<(), DispatchError> {
1447    let force_self_contained = retry == Some(RetryReason::LostDeltaBase);
1448    // Both wire names must fit SPEC-REFS §3's bound; a local branch named
1449    // before `MAX_BRANCH_NAME_BYTES` existed may not. Say so by name
1450    // rather than as a transport's "invalid ref name".
1451    refs::check_pushable_branch(branch)?;
1452    // Diff against the remote's CURRENT tip so we send only what it lacks
1453    // and can delta against bases it already holds. Planning is an
1454    // optimization; the head CAS below remains authoritative.
1455    // A caller that already planned against the remote's current tip (the
1456    // split decision does) passes its plan, so it is diffed once.
1457    let mut plan = if let Some(plan) = preplanned {
1458        plan
1459    } else {
1460        let remote_tip = tx.read_ref(&format!("refs/heads/{branch}"))?;
1461        // Plan FIRST, against the remote's current tip (#521 perf): a no-op push
1462        // (the remote already holds this closure) yields an empty plan and takes
1463        // the cheap head-only path below WITHOUT walking the packmap chain. Only
1464        // a push that actually has objects to send pays the O(depth) chain probe.
1465        transfer::plan_pack_with(
1466            store,
1467            tip,
1468            if force_self_contained {
1469                None
1470            } else {
1471                remote_tip
1472            },
1473            encode_delta_candidates_batch,
1474        )?
1475    };
1476
1477    let limits = tx.upload_limits();
1478    // Payload and serialized-byte limits are checked separately. Framing
1479    // overhead grows with the entry count, so a fixed margin is insufficient.
1480    let effective_cap = effective_payload_cap(pack_payload_cap, limits.max_pack_bytes)?;
1481
1482    if plan.is_empty() {
1483        // Nothing to send — the remote already holds the closure; just move
1484        // the head (no packmap change needed, so no chain walk and no atomic
1485        // two-ref advance).
1486        return commit_head(tx, &format!("refs/heads/{branch}"), condition, &tip, branch);
1487    }
1488
1489    if crate::signal::is_shutdown() {
1490        return Err(DispatchError::Interrupted);
1491    }
1492
1493    // Re-baseline decision (#406/#521), made only now that we know there IS
1494    // something to send: walk the current chain exactly once and decide
1495    // whether this push should reset the chain to a single self-contained
1496    // node instead of appending. When NOT re-baselining, the walk is cached
1497    // (`resolved_chain`) so `advance_packmap`'s append path can reuse it
1498    // instead of re-walking.
1499    //
1500    // A reset is committed together with the head; the ordered (non-atomic)
1501    // `advance_refs` fallback (packmap PUT then head PUT) would strand the
1502    // head at the old tip on a torn write, and a reset — unlike an append —
1503    // is not a superset of the prior chain, so it can't rebuild that stranded
1504    // closure. We therefore re-baseline ONLY when BOTH hold:
1505    //   * the transport advances both refs transactionally
1506    //     (`supports_atomic_advance`), AND
1507    //   * the head write is CAS-conditioned (not `Any`). A force push's `Any`
1508    //     head condition makes even an atomic transport fall back to the
1509    //     ordered two-PUT path (`Any` is not expressible on the atomic
1510    //     endpoint — see `HttpTransport::supports_atomic_advance`), so a
1511    //     force push MUST take the safe append path instead of resetting.
1512    let mut rebaseline = false;
1513    let mut resolved_chain = None;
1514    if !force_self_contained
1515        && rebaseline_threshold > 0
1516        && let Some(pm) = tx.read_ref(&packmap_ref(branch))?
1517    {
1518        match probe_chain(tx, branch, pm) {
1519            Ok(chain)
1520                if chain.depth + 1 > rebaseline_threshold
1521                    && tx.supports_atomic_advance()
1522                    && !matches!(condition, refs::RefWriteCondition::Any) =>
1523            {
1524                rebaseline = true;
1525            }
1526            Ok(chain) => resolved_chain = Some(chain),
1527            Err(DispatchError::PackChainInvalid { .. }) => {}
1528            Err(e) => return Err(e),
1529        }
1530    }
1531
1532    if rebaseline {
1533        // Force a full-closure plan: no external bases, so the pack is
1534        // self-contained and safe to reset the chain onto.
1535        let full = transfer::plan_pack_with(store, tip, None, encode_delta_candidates_batch)?;
1536        let full_cap = effective_payload_cap(pack_payload_cap, limits.max_pack_bytes)?;
1537        if estimated_ticketed_count(
1538            &estimate_pack_sizes(store, &full, full_cap, limits.max_pack_bytes)?,
1539            limits,
1540        ) > split::data_pack_budget(limits)
1541        {
1542            rebaseline = false;
1543        } else {
1544            plan = full;
1545        }
1546    }
1547
1548    // A retry re-planned the push: make sure the new plan still fits before
1549    // uploading any of it, so the failure leaves no half-uploaded advance.
1550    if let Some(reason) = retry
1551        && estimated_ticketed_count(
1552            &estimate_pack_sizes(store, &plan, effective_cap, limits.max_pack_bytes)?,
1553            limits,
1554        ) > split::data_pack_budget(limits)
1555        && let Err(DispatchError::PushTooLarge { packs, limit, .. }) =
1556            build_and_upload_packs(PackSink::Count, store, plan.clone(), effective_cap, limits)
1557    {
1558        return Err(DispatchError::RetryTooLarge {
1559            reason,
1560            packs,
1561            limit,
1562        });
1563    }
1564
1565    // No pre-flight PushTooLarge: the estimate uses uncompressed sizes, so it
1566    // would refuse pushes that compress into six packs. The seal-time gate is
1567    // exact and still stops before the seventh ticketed pack's BeginUpload.
1568    // (The re-baseline check above may use the conservative estimate: a false
1569    // "too many" only keeps the append plan.)
1570
1571    // Build the plan into one or more payload-bounded packs (splitting
1572    // when the plan exceeds `pack_payload_cap`, issue #831) and upload
1573    // each as it's sealed. Raws first (non-blobs before blobs), then
1574    // deltas (their bases are external — resolved from the fetcher's
1575    // store via earlier packs, never a base introduced earlier in THIS
1576    // push — so no in-pack ordering is required across the split,
1577    // SPEC-PACKFILE §4).
1578    //
1579    // `self_contained` is read out before `plan` moves into
1580    // `build_and_upload_packs` (which consumes it by value so its
1581    // delta streams can be handed to the compression fan-out without
1582    // an extra clone — see that function's doc comment).
1583    let self_contained = plan.self_contained;
1584    let head_ref = format!("refs/heads/{branch}");
1585    let pack_keys = build_and_upload_packs(
1586        PackSink::Upload {
1587            tx,
1588            head_ref: &head_ref,
1589        },
1590        store,
1591        plan,
1592        effective_cap,
1593        limits,
1594    )?;
1595
1596    // Chain the pack(s) onto the packmap AND move the head together
1597    // (#408): a transactional transport applies both atomically, the
1598    // default does packmap-then-head. Either way the head never lands
1599    // past a packmap that can't reconstruct it. `Append`'s
1600    // `self_contained` lets a full-closure push reset a broken chain
1601    // (unconditionally, on any transport); `ResetSelfContained`
1602    // proactively resets a healthy chain that has grown too deep, and
1603    // is only ever chosen above when the transport is atomic-capable
1604    // AND the head write is CAS-conditioned. A failed advance leaves
1605    // the head untouched.
1606    let action = if rebaseline {
1607        ChainAction::ResetSelfContained
1608    } else {
1609        ChainAction::Append { self_contained }
1610    };
1611    advance_packmap(
1612        tx,
1613        branch,
1614        &pack_keys,
1615        action,
1616        resolved_chain,
1617        condition,
1618        tip,
1619    )
1620}
1621
1622/// Build `plan`'s entries into one or more packs, each staying under
1623/// `payload_cap` bytes of wire payload, uploading each pack to `tx` as
1624/// soon as it's sealed. Returns the ordered pack keys — build order is
1625/// apply order, threaded straight into [`advance_packmap`].
1626///
1627/// A single linear left-to-right pass over the plan's already-ordered
1628/// `raw ++ deltas` sequence is enough: [`transfer::plan_pack`] never
1629/// deltas an entry against a base introduced earlier in THIS push
1630/// (every delta's base is already on the remote from a prior push), so
1631/// packs can be sealed independently the instant one would exceed the
1632/// cap — no intra-push base-ordering hazard to preserve across the
1633/// split.
1634///
1635/// Sizing uses a conservative *uncompressed* upper bound per entry
1636/// (`bytes.len()` for a raw, `HASH_LEN + stream.len()` for a delta)
1637/// checked against [`PackWriter::total_payload`]'s real (compressed)
1638/// running total, so a sealed pack never exceeds `payload_cap` — it may
1639/// under-fill when compression bites, yielding more packs than the
1640/// theoretical minimum, never fewer. Peak memory is one pack buffer
1641/// PLUS one prepared batch, not every pack (or the whole plan) held at
1642/// once: [`size_capped_batch_lens`] caps each batch at `payload_cap`
1643/// conservative bytes (the same budget the writer itself is bounded
1644/// by) as well as [`pack_fanout_threshold`]-many entries, so a plan
1645/// containing very large objects can't balloon a batch past what one
1646/// sealed pack would have held anyway — see that function's doc
1647/// comment for why both caps are needed.
1648///
1649/// Takes `plan` BY VALUE (not `&PackPlan`): the delta compression
1650/// fan-out needs owned `Vec<u8>` streams to hand to rayon, and
1651/// `plan.deltas` already owns them — draining them out here is zero
1652/// extra bytes copied, where borrowing would force a clone per entry
1653/// just to satisfy ownership. The caller reads whatever it needs off
1654/// `plan` (just `self_contained`) before making this call.
1655fn effective_payload_cap(requested: u64, advertised: Option<u64>) -> Result<u64, DispatchError> {
1656    // The serialized limit's overhead margin is exact and dynamic: header,
1657    // trailer and one frame per entry. `should_seal` applies it on each add.
1658    let cap = requested
1659        .min(pack::MAX_TOTAL_PAYLOAD)
1660        .min(advertised.unwrap_or(pack::MAX_TOTAL_PAYLOAD));
1661    if cap == 0 {
1662        return Err(DispatchError::Transport(TransportError::PayloadTooLarge(0)));
1663    }
1664    Ok(cap)
1665}
1666
1667fn serialized_bound(payload: u64, entries: usize) -> u64 {
1668    payload
1669        .saturating_add((entries as u64).saturating_mul(pack::ENTRY_FRAME_LEN as u64))
1670        .saturating_add((pack::HEADER_LEN + pack::TRAILER_LEN) as u64)
1671}
1672
1673fn ticketed_pack(limits: UploadLimits, serialized_size: u64) -> bool {
1674    limits.tickets_per_advance.is_some()
1675        && serialized_size >= limits.ticket_threshold_bytes.unwrap_or(0)
1676}
1677
1678fn estimated_ticketed_count(sizes: &[u64], limits: UploadLimits) -> usize {
1679    sizes
1680        .iter()
1681        .filter(|&&size| ticketed_pack(limits, size))
1682        .count()
1683}
1684
1685fn estimate_pack_sizes(
1686    store: &ObjectStore,
1687    plan: &transfer::PackPlan,
1688    payload_cap: u64,
1689    max_pack_bytes: Option<u64>,
1690) -> Result<Vec<u64>, DispatchError> {
1691    let mut packs = Vec::new();
1692    let mut payload = 0_u64;
1693    let mut entries = 0_usize;
1694    let sizes = plan
1695        .raw
1696        .iter()
1697        .map(|hash| store.object_metadata(hash).map(|meta| meta.len()))
1698        .collect::<Result<Vec<_>, _>>()?;
1699    for size in sizes.into_iter().chain(
1700        plan.deltas
1701            .iter()
1702            .map(|delta| (HASH_LEN + delta.stream.len()) as u64),
1703    ) {
1704        let next_payload = payload.saturating_add(size);
1705        if entries > 0
1706            && (next_payload > payload_cap
1707                || max_pack_bytes
1708                    .is_some_and(|limit| serialized_bound(next_payload, entries + 1) > limit))
1709        {
1710            packs.push(serialized_bound(payload, entries));
1711            payload = 0;
1712            entries = 0;
1713        }
1714        payload = payload.saturating_add(size);
1715        entries += 1;
1716    }
1717    if entries > 0 {
1718        packs.push(serialized_bound(payload, entries));
1719    }
1720    Ok(packs)
1721}
1722
1723/// Where sealed packs go: to the remote, or nowhere (an exact local dry run of
1724/// the seal, for the split's pre-flight check).
1725#[derive(Clone, Copy)]
1726enum PackSink<'a> {
1727    Upload {
1728        tx: &'a dyn Transport,
1729        head_ref: &'a str,
1730    },
1731    Count,
1732}
1733
1734fn build_and_upload_packs(
1735    sink: PackSink<'_>,
1736    store: &ObjectStore,
1737    plan: transfer::PackPlan,
1738    payload_cap: u64,
1739    limits: UploadLimits,
1740) -> Result<Vec<Hash>, DispatchError> {
1741    let mut pack_keys = Vec::new();
1742    let mut ticketed_count = 0;
1743    let mut w = PackWriter::new();
1744    let max_entries = pack_fanout_threshold().max(1);
1745    let transfer::PackPlan {
1746        raw, mut deltas, ..
1747    } = plan;
1748
1749    // Raw entries: size each via a metadata `stat` (cheap — no content
1750    // read) so batches can be capped by real bytes, not just count,
1751    // *before* any object's content is actually read into memory.
1752    let raw_sizes: Vec<u64> = raw
1753        .iter()
1754        .map(|h| Ok(store.object_metadata(h)?.len()))
1755        .collect::<Result<_, DispatchError>>()?;
1756    let mut start = 0;
1757    for len in size_capped_batch_lens(&raw_sizes, payload_cap, max_entries) {
1758        let chunk = &raw[start..start + len];
1759        start += len;
1760        for entry in prepare_raw_batch(store, chunk)? {
1761            if should_seal(
1762                &w,
1763                entry.conservative_len() as u64,
1764                payload_cap,
1765                limits.max_pack_bytes,
1766            ) {
1767                seal_pack(sink, &mut w, &mut pack_keys, &mut ticketed_count, limits)?;
1768            }
1769            w.push_prepared_raw(entry)?;
1770            // Honest progress (#711): one real object just got staged into
1771            // the outgoing pack. Never git's fabricated
1772            // Enumerating/Counting/Compressing lines — see `crate::progress`.
1773            report_packed(sink);
1774        }
1775    }
1776
1777    // Delta entries: sizes are already known (`stream.len()`, no I/O),
1778    // so batch lengths can be computed directly off them; `drain` then
1779    // moves each batch's `PlannedDelta`s out of `deltas` without
1780    // cloning their streams.
1781    let delta_sizes: Vec<u64> = deltas
1782        .iter()
1783        .map(|d| (HASH_LEN + d.stream.len()) as u64)
1784        .collect();
1785    for len in size_capped_batch_lens(&delta_sizes, payload_cap, max_entries) {
1786        let chunk: Vec<transfer::PlannedDelta> = deltas.drain(..len).collect();
1787        for entry in prepare_delta_batch(chunk) {
1788            if should_seal(
1789                &w,
1790                entry.conservative_len() as u64,
1791                payload_cap,
1792                limits.max_pack_bytes,
1793            ) {
1794                seal_pack(sink, &mut w, &mut pack_keys, &mut ticketed_count, limits)?;
1795            }
1796            w.push_prepared_delta(entry)?;
1797            report_packed(sink);
1798        }
1799    }
1800
1801    seal_pack(sink, &mut w, &mut pack_keys, &mut ticketed_count, limits)?;
1802    Ok(pack_keys)
1803}
1804
1805/// Honest progress (#711): one real object just got staged into the outgoing
1806/// pack. Never git's fabricated Enumerating/Counting/Compressing lines — see
1807/// `crate::progress`. A dry seal stages nothing to send and reports nothing.
1808fn report_packed(sink: PackSink<'_>) {
1809    if matches!(sink, PackSink::Upload { .. }) {
1810        crate::progress::report(crate::progress::Event::ObjectsPacked(1));
1811    }
1812}
1813
1814/// Split `sizes` (each item's conservative uncompressed byte size, in
1815/// original order) into batch lengths such that each batch's summed
1816/// size stays at-or-under `payload_cap` bytes AND at-or-under
1817/// `max_entries` items — whichever limit a batch hits first ends it.
1818/// Mirrors [`should_seal`]'s "never true for an empty writer" escape
1819/// hatch: a single item over `payload_cap` on its own (an oversized
1820/// blob against a tiny test cap, say) still gets a one-item batch of
1821/// its own rather than looping forever or being silently dropped.
1822///
1823/// The byte cap exists so [`build_and_upload_packs`]'s per-batch
1824/// compression fan-out can't hold dramatically more raw content in
1825/// memory at once than the single pack buffer it's feeding already
1826/// does — without it, a count-only cap (`max_entries`, sized for
1827/// rayon dispatch overhead, not memory) could group hundreds of
1828/// near-`MAX_RAW_OBJECT_SIZE` (1 GiB) blobs into one batch and read +
1829/// compress all of them before any of `payload_cap`'s own bookkeeping
1830/// ever runs. The entry cap exists for the opposite corpus shape: many
1831/// small objects would otherwise form one enormous `payload_cap`-sized
1832/// (default 4 GiB) batch, both blowing past `pack_fanout_threshold`'s
1833/// bench-measured parallel-dispatch sizing and stalling progress
1834/// reporting for the whole batch's compression time.
1835fn size_capped_batch_lens(sizes: &[u64], payload_cap: u64, max_entries: usize) -> Vec<usize> {
1836    let mut lens = Vec::new();
1837    let mut running: u64 = 0;
1838    let mut count = 0usize;
1839    for &size in sizes {
1840        if count > 0 && (running.saturating_add(size) > payload_cap || count >= max_entries) {
1841            lens.push(count);
1842            running = 0;
1843            count = 0;
1844        }
1845        running += size;
1846        count += 1;
1847    }
1848    if count > 0 {
1849        lens.push(count);
1850    }
1851    lens
1852}
1853
1854/// Entries-per-thread budget below which [`prepare_raw_batch`] /
1855/// [`prepare_delta_batch`] read + zstd-compress sequentially instead of
1856/// fanning out across rayon's thread pool, for a pool of a given size.
1857/// Same crossover shape as `add.rs`'s `HASH_FANOUT_FILES_PER_THREAD`
1858/// (PR #951) — rayon's pool-dispatch overhead loses to a plain loop
1859/// below a few entries per thread and wins clearly above it — measured
1860/// for the pack-build path specifically with `cargo bench -p
1861/// mkit-benches --bench pack_build_fanout`.
1862///
1863/// Also doubles as one of [`size_capped_batch_lens`]'s two caps: a
1864/// batch that isn't cut short by `payload_cap` first comes out exactly
1865/// this many entries long (until the final, possibly-shorter
1866/// remainder) — the same length `prepare_raw_batch`/
1867/// `prepare_delta_batch` compare against to decide sequential vs.
1868/// parallel, so a full-length batch always takes the parallel path.
1869const PACK_FANOUT_ENTRIES_PER_THREAD: usize = 8;
1870
1871/// The per-batch entry count [`build_and_upload_packs`] fans out across
1872/// rayon at — see [`PACK_FANOUT_ENTRIES_PER_THREAD`].
1873fn pack_fanout_threshold() -> usize {
1874    crate::fanout::threshold(PACK_FANOUT_ENTRIES_PER_THREAD)
1875}
1876
1877/// Read + compression-prepare one chunk of raw object hashes —
1878/// sequentially below [`pack_fanout_threshold`], via rayon's global
1879/// thread pool at or above it. Output mirrors `chunk`'s order 1:1
1880/// either way: `plan.raw` is already ordered (non-blobs before blobs,
1881/// each group in BLAKE3 order — see [`transfer::PackPlan`]'s doc
1882/// comment), and that order is the wire order the pack is built in, so
1883/// a parallel fan-out must not reshuffle it.
1884fn prepare_raw_batch(
1885    store: &ObjectStore,
1886    chunk: &[Hash],
1887) -> Result<Vec<PreparedRaw>, DispatchError> {
1888    crate::fanout::try_map_seq_or_par(chunk, pack_fanout_threshold(), |h| {
1889        Ok(PackWriter::prepare_raw(*h, store.read(h)?))
1890    })
1891}
1892
1893/// The delta-entry counterpart of [`prepare_raw_batch`]. Delta streams
1894/// are already resident in `chunk` (no store read needed), so this
1895/// only fans out the zstd compression step. Takes `chunk` by value —
1896/// `build_and_upload_packs` drains it straight out of `plan.deltas` —
1897/// so `PackWriter::prepare_delta` can take each stream by move instead
1898/// of cloning it.
1899fn prepare_delta_batch(chunk: Vec<transfer::PlannedDelta>) -> Vec<PreparedDelta> {
1900    if chunk.len() < pack_fanout_threshold() {
1901        return chunk
1902            .into_iter()
1903            .map(|d| PackWriter::prepare_delta(d.base, d.stream))
1904            .collect();
1905    }
1906    chunk
1907        .into_par_iter()
1908        .map(|d| PackWriter::prepare_delta(d.base, d.stream))
1909        .collect()
1910}
1911
1912/// Entries-per-thread budget below which [`encode_delta_candidates_batch`]
1913/// diffs candidates sequentially instead of fanning out across rayon's
1914/// thread pool, for a pool of a given size. Same crossover shape as
1915/// [`PACK_FANOUT_ENTRIES_PER_THREAD`] but tuned separately, and lower:
1916/// each entry here pays two disk reads (target + base) plus
1917/// `delta::encode`'s block-hash-table build and greedy scan over the
1918/// target's full length — a heavier per-item cost than the
1919/// zstd-compression fan-out `PACK_FANOUT_ENTRIES_PER_THREAD` guards, so
1920/// fewer entries per thread are enough to amortize dispatch.
1921///
1922/// `cargo bench -p mkit-benches --bench delta_plan_fanout` on a 4-core
1923/// host, isolating `delta::encode` over synthetic 64 KiB
1924/// (FastCDC-average-sized) near-duplicate chunks:
1925///
1926///   n=256: sequential 46.99ms -> rayon 14.79ms (3.18x)
1927///   n=128: sequential 22.85ms -> rayon  7.57ms (3.02x)
1928///   n=64:  sequential 11.65ms -> rayon  3.67ms (3.18x)
1929///   n=32:  sequential  5.52ms -> rayon  1.74ms (3.18x)
1930///   n=16:  sequential  3.28ms -> rayon  1.22ms (2.69x)
1931///   n=8:   sequential  1.38ms -> rayon  0.70ms (1.97x)
1932///   n=4:   sequential  0.68ms -> rayon  0.42ms (1.63x)
1933///   n=2:   sequential  0.34ms -> rayon  0.33ms (roughly flat)
1934///   n=1:   sequential  0.17ms -> rayon  0.17ms (roughly flat)
1935///
1936/// Each candidate already costs ~0.17ms, comfortably above rayon's
1937/// microsecond-scale dispatch overhead, so the crossover falls at the
1938/// smallest candidate count worth fanning out at all rather than at a
1939/// larger multiple of the pool size the other fan-outs need.
1940const DELTA_PLAN_FANOUT_ENTRIES_PER_THREAD: usize = 1;
1941
1942/// The candidate count [`encode_delta_candidates_batch`] fans out across
1943/// rayon at — see [`DELTA_PLAN_FANOUT_ENTRIES_PER_THREAD`].
1944fn delta_plan_fanout_threshold() -> usize {
1945    crate::fanout::threshold(DELTA_PLAN_FANOUT_ENTRIES_PER_THREAD)
1946}
1947
1948/// [`transfer::plan_pack_with`]'s `encode_deltas` callback: diffs every
1949/// delta candidate `plan_pack_with` found — sequentially below
1950/// [`delta_plan_fanout_threshold`], via rayon's global thread pool at or
1951/// above it. Each candidate's `delta::encode` is independent of every
1952/// other candidate's (same shape as `prepare_raw_batch`'s per-object
1953/// compression fan-out), so a push touching many changed blobs
1954/// parallelizes the diffing itself, not just the downstream zstd pass
1955/// `prepare_delta_batch` already fans out.
1956///
1957/// Processes at most [`DELTA_CANDIDATE_BATCH_CAP`] candidates at a time.
1958/// Repeated bases share a cache capped at [`DELTA_BASE_CACHE_BYTES`];
1959/// one-use bases and bases that do not fit are read by each encoder as
1960/// before. Metadata checks avoid reading nonfitting bases twice. The
1961/// cache is dropped before the next batch. Active encoders' target/base
1962/// buffers and the accumulated encoded output are additional allocations;
1963/// imported blobs can reach the store's 1 GiB object limit.
1964///
1965/// Order matches `candidates` 1:1, including across batch boundaries,
1966/// satisfying `plan_pack_with`'s contract.
1967fn encode_delta_candidates_batch(
1968    store: &ObjectStore,
1969    candidates: &[transfer::DeltaCandidate],
1970) -> Result<Vec<Option<transfer::PlannedDelta>>, StoreError> {
1971    let mut results = Vec::with_capacity(candidates.len());
1972    for batch in candidates.chunks(DELTA_CANDIDATE_BATCH_CAP) {
1973        let base_bytes = cache_delta_bases(store, batch, DELTA_BASE_CACHE_BYTES)?;
1974        results.extend(crate::fanout::try_map_seq_or_par(
1975            batch,
1976            delta_plan_fanout_threshold(),
1977            |&c| match base_bytes.get(&c.base) {
1978                Some(bytes) => transfer::encode_delta_candidate_with_base(store, c, bytes),
1979                None => transfer::encode_delta_candidate(store, c),
1980            },
1981        )?);
1982    }
1983    Ok(results)
1984}
1985
1986/// Maximum number of candidates whose base objects are cached together.
1987const DELTA_CANDIDATE_BATCH_CAP: usize = 64;
1988
1989/// Maximum retained payload bytes for reused bases, independent of blob sizes.
1990const DELTA_BASE_CACHE_BYTES: usize = 64 * 1024 * 1024;
1991
1992fn cache_delta_bases(
1993    store: &ObjectStore,
1994    candidates: &[transfer::DeltaCandidate],
1995    byte_budget: usize,
1996) -> Result<std::collections::HashMap<Hash, Vec<u8>>, StoreError> {
1997    let mut counts = std::collections::BTreeMap::<Hash, usize>::new();
1998    for candidate in candidates {
1999        *counts.entry(candidate.base).or_default() += 1;
2000    }
2001    let mut cache = std::collections::HashMap::new();
2002    let mut remaining = byte_budget;
2003    for (base, count) in counts {
2004        if count < 2 || remaining == 0 {
2005            continue;
2006        }
2007        if store.object_metadata(&base)?.len() > remaining as u64 {
2008            continue;
2009        }
2010        let bytes = store.read(&base)?;
2011        // Recheck after reading: metadata is only an optimization, never
2012        // authority for the retained byte bound if a file changes meanwhile.
2013        if bytes.len() <= remaining {
2014            remaining -= bytes.len();
2015            cache.insert(base, bytes);
2016        }
2017    }
2018    Ok(cache)
2019}
2020
2021/// Would pushing an entry of (conservative, uncompressed) size
2022/// `add_bound` into `w` exceed `payload_cap`? Never true for an empty
2023/// writer — a single entry over the cap (only reachable with a
2024/// test-injected tiny cap; production entries are bounded well under
2025/// [`pack::MAX_TOTAL_PAYLOAD`] by [`mkit_core::store::MAX_RAW_OBJECT_SIZE`])
2026/// lands alone in its own pack rather than looping forever.
2027fn should_seal(
2028    w: &PackWriter,
2029    add_bound: u64,
2030    payload_cap: u64,
2031    max_pack_bytes: Option<u64>,
2032) -> bool {
2033    let next_payload = w.total_payload().saturating_add(add_bound);
2034    w.entry_count() > 0
2035        && (next_payload > payload_cap
2036            || max_pack_bytes
2037                .is_some_and(|limit| serialized_bound(next_payload, w.entry_count() + 1) > limit))
2038}
2039
2040/// Finish `w`, upload it, record its key, and replace `w` with a fresh
2041/// empty writer so the caller can keep pushing entries into the next
2042/// pack.
2043fn seal_pack(
2044    sink: PackSink<'_>,
2045    w: &mut PackWriter,
2046    pack_keys: &mut Vec<Hash>,
2047    ticketed_count: &mut usize,
2048    limits: UploadLimits,
2049) -> Result<(), DispatchError> {
2050    if crate::signal::is_shutdown() {
2051        return Err(DispatchError::Interrupted);
2052    }
2053    let sealed = std::mem::replace(w, PackWriter::new());
2054    let pack = sealed.finish()?;
2055    let is_ticketed = ticketed_pack(limits, pack.len() as u64);
2056    let budget = split::data_pack_budget(limits);
2057    if is_ticketed && *ticketed_count >= budget {
2058        return Err(DispatchError::PushTooLarge {
2059            packs: *ticketed_count + 1,
2060            limit: budget,
2061            commit: None,
2062            holds_remote_head: false,
2063        });
2064    }
2065    if limits
2066        .max_pack_bytes
2067        .is_some_and(|limit| pack.len() as u64 > limit)
2068    {
2069        return Err(DispatchError::Transport(TransportError::PayloadTooLarge(
2070            pack.len(),
2071        )));
2072    }
2073    let pack_key = pack::pack_key(&pack);
2074    if let PackSink::Upload { tx, head_ref } = sink {
2075        tx.upload_pack_via_ref(&pack, &PackKey::from_hash(pack_key), head_ref)
2076            .map_err(|error| repository_operation_error(tx, error))?;
2077        // Upload is complete — report the real byte count handed to the
2078        // transport, not an estimate.
2079        crate::progress::report(crate::progress::Event::PackUploaded(pack.len() as u64));
2080    }
2081    pack_keys.push(pack_key);
2082    if is_ticketed {
2083        *ticketed_count += 1;
2084    }
2085    Ok(())
2086}
2087
2088/// [`pull_all_with`] with signature verification on — the CLI's default
2089/// (issue #692). Existing in-process callers (the integration-test suite)
2090/// that construct only validly-signed histories are unaffected.
2091pub fn pull_all(
2092    cwd: &Path,
2093    tx: &dyn Transport,
2094    remote: &str,
2095    target_branch: Option<&str>,
2096) -> Result<usize, DispatchError> {
2097    pull_all_with(cwd, tx, remote, target_branch, true)
2098}
2099
2100/// Fetch remote refs, then fast-forward the current local branch from
2101/// `refs/remotes/default/<branch>`. Fresh repos with no local branch tip
2102/// initialise from the current branch's remote-tracking ref, or the first
2103/// advertised remote branch when the current default branch is absent.
2104///
2105/// `target_branch`, when `Some`, overrides which remote branch to land
2106/// on (used by `mkit clone -b <branch>`): the branch MUST exist among
2107/// the remote's advertised refs or the call fails with
2108/// [`DispatchError::RemoteBranchMissing`] rather than silently falling
2109/// back to another branch. `None` preserves the historical HEAD-driven
2110/// selection used by plain `pull`.
2111///
2112/// `require_signed` gates the post-fetch commit/remix/tag signature check
2113/// (issue #692) — `true` (the CLI's default, see [`pull_all`]) verifies
2114/// every newly-fetched object and fails closed; `false` is the explicit
2115/// `--no-verify-signatures` / `pull.require_signed = false` opt-out.
2116pub fn pull_all_with(
2117    cwd: &Path,
2118    tx: &dyn Transport,
2119    remote: &str,
2120    target_branch: Option<&str>,
2121    require_signed: bool,
2122) -> Result<usize, DispatchError> {
2123    let layout = mkit_core::layout::discover(cwd)?;
2124    let store = crate::commands::open_store_configured(&layout)?;
2125    // Fetch phase: `fetch_objects` takes the repo lock itself, narrowly and
2126    // per branch, around only the local unpack + remote-ref-publish window
2127    // (#642 — see `packmap::apply_fetched_chain`). No lock is held here
2128    // across the network transfer.
2129    let n = fetch_objects(&store, &layout, tx, remote, target_branch, require_signed)?;
2130    let remote_refs = crate::commands::list_remote_refs_parallel(&layout, remote)?
2131        .into_iter()
2132        .filter_map(|r| r.hash.map(|hash| (r.name, hash)))
2133        .collect::<Vec<_>>();
2134    if remote_refs.is_empty() {
2135        return Ok(n);
2136    }
2137
2138    // Fast-forward phase (#642): branch ref + HEAD + worktree, narrowly
2139    // locked around just this window rather than bundled with the fetch
2140    // above. The objects `remote_tip` points at are already reachable via
2141    // the remote-tracking ref published during the fetch phase, so this
2142    // lock's job here is worktree-mutation exclusivity against concurrent
2143    // commands (e.g. a racing `commit`/`checkout`/`reset`), not GC
2144    // protection — that hazard was already closed before this lock was
2145    // taken.
2146    let _lock = mkit_core::repo_lock::acquire_default(
2147        layout.worktree_state_dir(),
2148        crate::commands::WORKTREE_LOCK,
2149    )?;
2150    crate::commands::warn_if_served(&layout);
2151    let original_head = refs::read_head(&layout).ok();
2152    let (branch, local_tip, remote_tip) = match &original_head {
2153        Some(Head::Branch(head_branch)) => {
2154            let want_branch = target_branch.unwrap_or(head_branch.as_str());
2155            let local_tip = refs::read_ref(&layout, want_branch)?;
2156            let selected = if local_tip.is_some() || target_branch.is_some() {
2157                // An explicit `-b <branch>` (or an already-committed local
2158                // branch of that name) must match exactly — no silent
2159                // fallback to a different branch.
2160                remote_refs
2161                    .iter()
2162                    .find(|(name, _)| name == want_branch)
2163                    .ok_or_else(|| DispatchError::RemoteBranchMissing(want_branch.to_owned()))?
2164            } else {
2165                remote_refs
2166                    .iter()
2167                    .find(|(name, _)| name == want_branch)
2168                    .unwrap_or(&remote_refs[0])
2169            };
2170            (selected.0.clone(), local_tip, selected.1)
2171        }
2172        Some(Head::Detached(_)) => return Err(DispatchError::DetachedHead),
2173        None => (remote_refs[0].0.clone(), None, remote_refs[0].1),
2174    };
2175
2176    let ref_condition = if let Some(local_tip) = local_tip {
2177        if local_tip == remote_tip {
2178            return Ok(n);
2179        }
2180        if !is_ancestor(&store, local_tip, remote_tip)? {
2181            return Err(DispatchError::NonFastForwardPull { branch });
2182        }
2183        refs::RefWriteCondition::Match(local_tip)
2184    } else {
2185        refs::RefWriteCondition::Missing
2186    };
2187
2188    let tree = load_tree_hash(&store, remote_tip)?;
2189    crate::commands::ensure_restore_safe(&layout, &store, tree)
2190        .map_err(DispatchError::RestoreSafety)?;
2191    crate::commands::write_ref_recording_history(&layout, &branch, ref_condition, &remote_tip)?;
2192    if let Err(e) = refs::write_head_branch(&layout, &branch) {
2193        rollback_pull_ref(&layout, &branch, local_tip, remote_tip)?;
2194        return Err(e.into());
2195    }
2196    if let Err(e) = crate::commands::restore_worktree_and_index(&layout, &store, tree) {
2197        if let Err(rollback) =
2198            rollback_pull_ref_and_head(&layout, &branch, local_tip, remote_tip, original_head)
2199        {
2200            return Err(DispatchError::RestoreSafety(format!(
2201                "{e}; additionally failed to roll back ref: {rollback}"
2202            )));
2203        }
2204        return Err(DispatchError::RestoreSafety(e));
2205    }
2206    Ok(n)
2207}
2208
2209fn rollback_pull_ref_and_head(
2210    layout: &RepoLayout,
2211    branch: &str,
2212    local_tip: Option<Hash>,
2213    remote_tip: Hash,
2214    original_head: Option<Head>,
2215) -> Result<(), String> {
2216    rollback_pull_ref(layout, branch, local_tip, remote_tip).map_err(|e| e.to_string())?;
2217    match original_head {
2218        Some(Head::Branch(name)) => refs::write_head_branch(layout, &name),
2219        Some(Head::Detached(hash)) => refs::write_head_detached(layout, &hash),
2220        None => Ok(()),
2221    }
2222    .map_err(|e| e.to_string())
2223}
2224
2225fn rollback_pull_ref(
2226    layout: &RepoLayout,
2227    branch: &str,
2228    local_tip: Option<Hash>,
2229    remote_tip: Hash,
2230) -> Result<(), refs::RefError> {
2231    if let Some(local_tip) = local_tip {
2232        crate::commands::write_ref_recording_history(
2233            layout,
2234            branch,
2235            refs::RefWriteCondition::Match(remote_tip),
2236            &local_tip,
2237        )
2238    } else if refs::read_ref(layout, branch)? == Some(remote_tip) {
2239        refs::delete_ref(layout, branch)
2240    } else {
2241        Ok(())
2242    }
2243}
2244
2245/// [`fetch_all_with`] with signature verification on — the CLI's default
2246/// (issue #692). Existing in-process callers (the integration-test suite)
2247/// that construct only validly-signed histories are unaffected.
2248pub fn fetch_all(cwd: &Path, tx: &dyn Transport, remote: &str) -> Result<usize, DispatchError> {
2249    fetch_all_with(cwd, tx, remote, true)
2250}
2251
2252/// `fetch` — `pull_all` without the HEAD update. Downloads every object
2253/// reachable from each remote ref (via [`Transport::download_pack`] on
2254/// the object's own digest) and writes the ref into
2255/// `refs/remotes/default/<branch>`.
2256///
2257/// See [`pull_all_with`] for the `require_signed` contract (issue #692).
2258pub fn fetch_all_with(
2259    cwd: &Path,
2260    tx: &dyn Transport,
2261    remote: &str,
2262    require_signed: bool,
2263) -> Result<usize, DispatchError> {
2264    let layout = mkit_core::layout::discover(cwd)?;
2265    // No outer lock here (#642): `fetch_objects` takes the repo lock
2266    // itself, narrowly and per branch, around only the local unpack +
2267    // remote-ref-publish window for that branch — never across the
2268    // network transfer. See `packmap::resolve_and_download_chain` /
2269    // `apply_fetched_chain` and `fetch_objects_inner` below.
2270    let store = crate::commands::open_store_configured(&layout)?;
2271    fetch_objects(&store, &layout, tx, remote, None, require_signed)
2272}
2273
2274/// Reconstruct every remote `refs/heads/*` from its packmap chain and
2275/// publish the remote-tracking refs. Each branch's local object writes and
2276/// its ref publish happen under a repo lock acquired fresh for that branch
2277/// (#642) — see [`fetch_objects_inner`] — so the caller does NOT need to
2278/// hold the repo lock around this call.
2279///
2280/// mkit speaks a single, packmap-only transfer dialect (the legacy
2281/// per-object download path was removed): for every advertised branch the
2282/// flow is exactly
2283///
2284/// 1. read the branch's packmap ref (`refs/mkit/packmap/<branch>`),
2285/// 2. walk its chain oldest-first and download any pack the local
2286///    applied-pack record doesn't already have
2287///    ([`packmap::resolve_and_download_chain`]) — no repo lock held,
2288/// 3. acquire the repo lock, unpack the downloaded packs
2289///    ([`packmap::apply_fetched_chain`]), then
2290/// 4. assert the tip's closure is fully present
2291///    ([`verify_closure_present`]) — a pure integrity check that downloads
2292///    nothing (still under the lock from step 3), then publish the
2293///    branch's remote-tracking ref and release the lock (#642).
2294///
2295/// Steps 3 and 4 are both performed *inside* [`packmap::apply_fetched_chain`]
2296/// (rather than sequenced here) so a closure-completeness failure counts
2297/// toward that function's applied-pack self-heal retry (#409): if the
2298/// local record wrongly claims every pack in the chain is already applied
2299/// (e.g. `.mkit/objects` was wiped out-of-band while `applied-packs/`
2300/// survived), the very first symptom is exactly this closure check
2301/// failing, not a download/unpack error — the retry has to cover both.
2302///
2303/// Both ends fail loudly: an absent packmap beside a present head is
2304/// [`DispatchError::PackmapMissing`] (a head absent too is a stale listing,
2305/// STC §7.9: the branch is skipped and its tracking ref left untouched)
2306/// and a present-but-incomplete packmap (even after the self-heal retry) is
2307/// [`DispatchError::RemoteMissingObject`]. We never publish a
2308/// remote-tracking ref to a closure we couldn't fully materialise locally.
2309///
2310/// # Concurrent re-baseline (mkit #521)
2311///
2312/// We list branch tips (`list_refs`) BEFORE reading each branch's packmap,
2313/// so a concurrent push that re-baselines (resets the packmap to a fresh
2314/// single node, #406) between those two reads can leave us verifying an
2315/// OLD tip `h` against a packmap whose closure only covers the NEW tip —
2316/// surfacing as [`DispatchError::RemoteMissingObject`] (the reset chain
2317/// isn't a superset of the prior one, and the applied-pack self-heal can't
2318/// rescue a genuinely stale tip). This is transient — no bad ref was
2319/// published — so on that specific error we re-read the branch's CURRENT tip
2320/// and packmap and retry the chain once with the fresh pair, publishing the
2321/// fresh tip. A second failure (or a vanished branch) propagates unchanged.
2322///
2323/// # Applied-packs record: load once, persist once (mkit #546)
2324///
2325/// The applied-pack record (`<common dir>/applied-packs/<remote>`, #409) is
2326/// keyed by remote, not by branch, so this function — not
2327/// [`packmap::resolve_and_download_chain`] / [`packmap::apply_fetched_chain`]
2328/// — owns its lifecycle for the WHOLE fetch: loaded once before the branch
2329/// loop, persisted once after, however many branches are fetched. Those two
2330/// functions only mutate the record in memory (inserting applied digests,
2331/// or clearing on self-heal); neither touches disk. The final persist is
2332/// unconditional and best-effort — it runs on every outcome so a fetch that
2333/// applied packs before failing never loses that progress, and the record
2334/// is a pure performance cache whose own I/O must never fail a fetch whose
2335/// objects durably landed.
2336///
2337/// # Repo-lock scope (mkit #642)
2338///
2339/// Each branch acquires the repo lock fresh, right before
2340/// [`packmap::apply_fetched_chain`] unpacks that branch's downloaded packs,
2341/// and releases it right after the branch's remote-tracking ref is
2342/// published — see [`fetch_objects_inner`]. Nothing here holds the lock
2343/// during [`packmap::resolve_and_download_chain`]'s network I/O, and the
2344/// lock is released between branches, so only the local write + ref-publish
2345/// window for one branch at a time is ever locked. This still closes the
2346/// #267 GC-prune race: `gc` takes the very same lock before computing its
2347/// live set, so it can never observe a branch's objects on disk without
2348/// that branch's ref already published.
2349///
2350/// `target_branch` (a `clone -b` / `pull` target) is fetched even when an
2351/// eventual `ListRefs` (D34, STC §7.9) has not caught up with it: the branch
2352/// is strongly read by name, fetched if present, and
2353/// [`DispatchError::RemoteBranchMissing`] if absent. Never an exit-0 empty
2354/// clone for `-b`.
2355fn fetch_objects(
2356    store: &ObjectStore,
2357    layout: &RepoLayout,
2358    tx: &dyn Transport,
2359    remote: &str,
2360    target_branch: Option<&str>,
2361    require_signed: bool,
2362) -> Result<usize, DispatchError> {
2363    let mut applied = AppliedPacks::load_or_empty(layout, remote);
2364    let result = fetch_objects_inner(
2365        store,
2366        layout,
2367        tx,
2368        remote,
2369        target_branch,
2370        &mut applied,
2371        require_signed,
2372    );
2373    persist_record(&mut applied, remote);
2374    result
2375}
2376
2377/// The branch loop proper — see [`fetch_objects`] for the load-once /
2378/// persist-once applied-packs contract this is called under.
2379fn fetch_objects_inner(
2380    store: &ObjectStore,
2381    layout: &RepoLayout,
2382    tx: &dyn Transport,
2383    remote: &str,
2384    target_branch: Option<&str>,
2385    applied: &mut AppliedPacks,
2386    require_signed: bool,
2387) -> Result<usize, DispatchError> {
2388    let mut remote_refs = tx
2389        .list_refs("refs/heads/")
2390        .map_err(|error| repository_operation_error(tx, error))?;
2391    if let Some(branch) = target_branch
2392        && !remote_refs.iter().any(|r| r.name == branch)
2393    {
2394        // The listing may be stale (eventual under D34): a strong read by
2395        // name decides whether the branch exists. A transport error here
2396        // propagates; only `Ok(None)` is the "no such branch" verdict.
2397        match tx.read_ref(&format!("refs/heads/{branch}"))? {
2398            Some(hash) => remote_refs.push(refs::Ref {
2399                name: branch.to_owned(),
2400                hash: Some(hash),
2401            }),
2402            None => return Err(DispatchError::RemoteBranchMissing(branch.to_owned())),
2403        }
2404    }
2405    let mut n = 0;
2406    // Batch every fetched branch's remote-tracking-ref write (#645): see
2407    // `push_all_with` for the same pattern and its rationale. `tracking.write`
2408    // still runs inside this branch's `_lock` scope below, so ref visibility
2409    // timing is unchanged from the per-branch-publish-then-unlock model
2410    // (#642) — only the parent-directory fsync is deferred to one pass after
2411    // the loop. GC's #267 protection is unaffected: it depends on the
2412    // repo lock covering the object-write-to-ref-publish window (still true
2413    // per branch below), not on when the ref directory itself is fsynced.
2414    let mut tracking = refs::RemoteRefBatch::new(layout, remote)?;
2415    let result: Result<(), DispatchError> = (|| {
2416        for r in remote_refs {
2417            if crate::signal::is_shutdown() {
2418                return Err(DispatchError::Interrupted);
2419            }
2420            let Some(h) = r.hash else { continue };
2421            // A branch tip without a packmap is a corrupt/incomplete remote, not
2422            // a format we degrade to: the push path ALWAYS advertises a packmap
2423            // before moving the branch ref. A real transport error (network blip,
2424            // auth) propagates unchanged — only `Ok(None)` is the explicit
2425            // "no packmap" verdict.
2426            //
2427            // The one exception is a stale listing (STC §7.9): `ListRefs` is
2428            // eventual, so under D34 a branch deleted since the listing can
2429            // still be named. Re-read the head strongly: absent too means the
2430            // listing is stale, so skip the branch and write no tracking ref;
2431            // present means corruption (a head and its packmap share a shard),
2432            // which stays `PackmapMissing`. A transport error on the re-read
2433            // propagates and never becomes a skip.
2434            let Some(chain_head) = tx.read_ref(&packmap_ref(&r.name))? else {
2435                if tx.read_ref(&format!("refs/heads/{}", r.name))?.is_none() {
2436                    if crate::progress::should_report(false) {
2437                        eprintln!(
2438                            "fetch: skipping branch `{}`: listed by an eventual ListRefs but \
2439                             no longer present (deleted since the listing)",
2440                            r.name
2441                        );
2442                    }
2443                    continue;
2444                }
2445                return Err(DispatchError::PackmapMissing(r.name.clone()));
2446            };
2447
2448            // Phase 1 (#642): resolve this branch's chain and download its
2449            // packs — pure network I/O, no repo lock held (see
2450            // `packmap::resolve_and_download_chain`).
2451            let fetched = resolve_and_download_chain(tx, &r.name, chain_head, applied)?;
2452
2453            // Phase 2 (#642): unpack + verify + publish this branch's ref,
2454            // under a repo lock acquired fresh for this branch and released
2455            // once this loop iteration ends — never held across another
2456            // branch's download. See `packmap::apply_fetched_chain`'s doc
2457            // comment for the safety contract (closing the #267 GC-prune
2458            // race) this depends on.
2459            let lock = mkit_core::repo_lock::acquire_default(
2460                layout.worktree_state_dir(),
2461                crate::commands::WORKTREE_LOCK,
2462            )?;
2463            crate::commands::warn_if_served(layout);
2464            // The tip we publish: normally the listed `h`, but if the chain fails
2465            // because a concurrent re-baseline moved the branch under us, the
2466            // freshly re-read tip (see this fn's doc comment). The match also
2467            // carries the lock guard out, so whichever branch runs, the
2468            // correct (possibly re-acquired) guard is what's held at
2469            // `tracking.write` below — see the retry branch's comment for why
2470            // there are two guards, not one held across the whole match.
2471            let (published_tip, _lock) = match apply_fetched_chain(
2472                store,
2473                tx,
2474                remote,
2475                &r.name,
2476                fetched,
2477                h,
2478                applied,
2479                require_signed,
2480            ) {
2481                Ok(()) => (h, lock),
2482                Err(e @ DispatchError::RemoteMissingObject(_)) => {
2483                    // Re-read the branch's CURRENT tip + packmap. If either
2484                    // is gone (branch deleted mid-fetch) the original error
2485                    // stands. Otherwise retry the chain ONCE with the fresh
2486                    // pair; a second failure propagates via `?`.
2487                    let (Some(fresh_h), Some(fresh_head)) = (
2488                        tx.read_ref(&format!("refs/heads/{}", r.name))?,
2489                        tx.read_ref(&packmap_ref(&r.name))?,
2490                    ) else {
2491                        return Err(e);
2492                    };
2493                    // Release the lock for the retry's network
2494                    // re-download too — mirrors phase 1's unlocked
2495                    // download exactly, rather than the previously
2496                    // "accepted trade" of holding the lock across a
2497                    // second network round-trip on this rare
2498                    // race-recovery path. Re-acquire before the
2499                    // retry's local unpack + publish, which still
2500                    // needs the same #267 protection phase 2 always
2501                    // has.
2502                    drop(lock);
2503                    let fresh_fetched =
2504                        resolve_and_download_chain(tx, &r.name, fresh_head, applied)?;
2505                    let lock = mkit_core::repo_lock::acquire_default(
2506                        layout.worktree_state_dir(),
2507                        crate::commands::WORKTREE_LOCK,
2508                    )?;
2509                    crate::commands::warn_if_served(layout);
2510                    apply_fetched_chain(
2511                        store,
2512                        tx,
2513                        remote,
2514                        &r.name,
2515                        fresh_fetched,
2516                        fresh_h,
2517                        applied,
2518                        require_signed,
2519                    )?;
2520                    (fresh_h, lock)
2521                }
2522                Err(e) => return Err(e),
2523            };
2524            // Still inside `_lock`'s scope (the original guard, or the
2525            // retry's re-acquired one — see above).
2526            tracking.write(&r.name, &published_tip)?;
2527            n += 1;
2528        }
2529        Ok(())
2530    })();
2531    // Commit whatever tracking-ref writes succeeded regardless of how the
2532    // loop above ended — a mid-loop failure still durably publishes the
2533    // prefix that already verified successfully, matching the old
2534    // per-branch loop's per-write durability.
2535    tracking.commit()?;
2536    result?;
2537    Ok(n)
2538}
2539
2540/// Best-effort persist of `applied` (a write failure is logged and
2541/// swallowed), called exactly once per fetch — see [`fetch_objects`] for
2542/// the load-once / persist-once contract.
2543fn persist_record(applied: &mut AppliedPacks, remote: &str) {
2544    if let Err(e) = applied.persist() {
2545        eprintln!(
2546            "warning: could not persist applied-packs record for remote '{remote}' ({e}); it will be rebuilt on the next fetch"
2547        );
2548    }
2549}
2550
2551/// Assert that every object reachable from `tip` is already present in the
2552/// local store after the packmap chain has been unpacked. This is a pure
2553/// integrity check — it walks the closure via
2554/// [`mkit_core::ops::reachable_closure_checked`] (which reads each object and
2555/// re-verifies its digest) and performs NO network access. A reachable object
2556/// that the chain failed to deliver surfaces as [`StoreError::ObjectNotFound`],
2557/// which we re-tag as [`DispatchError::RemoteMissingObject`] so the fetch
2558/// aborts loudly rather than publishing a ref to a closure we can't
2559/// reconstruct.
2560///
2561/// When packs were skipped (the applied-pack fast path, #409) this walk is
2562/// the *sole* guarantee the store is complete, so it must not silently pass
2563/// on an unverified frontier: a closure exceeding the
2564/// [`mkit_core::ops::graph::MAX_REACHABLE`] cap leaves objects past the cap
2565/// unchecked, which over a partially-wiped store could hide missing objects.
2566/// We therefore surface truncation as a hard [`DispatchError::ClosureTooLarge`]
2567/// rather than dropping the flag. `ClosureTooLarge` is deliberately distinct
2568/// from `RemoteMissingObject` so it does NOT feed the self-heal retry — a
2569/// too-large history is not evidence of local staleness.
2570///
2571/// Called from [`packmap::fetch_pack_chain`] (not sequenced after it) so its
2572/// `RemoteMissingObject` result participates in that function's applied-pack
2573/// self-heal retry — see [`fetch_objects`]'s doc comment.
2574pub(crate) fn verify_closure_present(store: &ObjectStore, tip: &Hash) -> Result<(), DispatchError> {
2575    match mkit_core::ops::reachable_closure_checked(store, std::iter::once(tip)) {
2576        Ok((_, false)) => Ok(()),
2577        Ok((_, true)) => Err(DispatchError::ClosureTooLarge(
2578            mkit_core::ops::graph::MAX_REACHABLE,
2579        )),
2580        Err(StoreError::ObjectNotFound(hex)) => Err(DispatchError::RemoteMissingObject(hex)),
2581        Err(e) => Err(e.into()),
2582    }
2583}
2584
2585fn load_tree_hash(store: &ObjectStore, commit_hash: Hash) -> Result<Hash, DispatchError> {
2586    match store.read_object(&commit_hash)? {
2587        Object::Commit(c) => Ok(c.tree_hash),
2588        Object::Remix(r) => Ok(r.tree_hash),
2589        _ => Err(DispatchError::NotCommit),
2590    }
2591}
2592
2593#[cfg(test)]
2594mod tests {
2595    use super::ssh_options_from_config;
2596    use crate::config::Config;
2597    use mkit_core::layout::RepoLayout;
2598    use mkit_core::pack::{PreparedDelta, PreparedRaw};
2599    use mkit_core::store::ObjectStore;
2600    use mkit_core::transfer;
2601
2602    // =================================================================
2603    // `size_capped_batch_lens` — the batch-boundary arithmetic
2604    // `build_and_upload_packs` uses to bound each compression-fan-out
2605    // batch by both bytes (`payload_cap`) and entry count
2606    // (`pack_fanout_threshold`). Pure and deterministic (no rayon pool
2607    // size involved), so these pin the boundary conditions directly
2608    // rather than relying on an end-to-end push to happen to exercise
2609    // them on whatever hardware runs the test.
2610    // =================================================================
2611
2612    #[test]
2613    fn envelope_destination_trust_precedes_missing_key_resolution() {
2614        let directory = tempfile::tempdir().unwrap();
2615        let layout = RepoLayout::single(directory.path());
2616        for endpoint in ["mkit+https://untrusted.example", "mkit+http://127.0.0.1:1"] {
2617            let mut cfg = Config {
2618                transport_auth: "envelope".into(),
2619                signer: "legacy".into(),
2620                signing_key: ".mkit/keys/missing-signing-key".into(),
2621                ..Config::default()
2622            };
2623            let error = super::open_with_config(endpoint, &cfg, &layout)
2624                .err()
2625                .expect("untrusted signing must fail");
2626            assert!(
2627                matches!(error, super::DispatchError::UntrustedRemote(_)),
2628                "destination rejection must precede key access: {error}"
2629            );
2630            cfg.trusted_remote_endpoint = endpoint.into();
2631            let error = super::open_with_config(endpoint, &cfg, &layout)
2632                .err()
2633                .expect("trusted endpoint now resolves the absent key");
2634            assert!(
2635                error.to_string().contains("requires a signing key"),
2636                "trusted destination should reach key resolution: {error}"
2637            );
2638        }
2639    }
2640
2641    #[test]
2642    fn configured_connect_open_reports_malformed_repository_identity() {
2643        let directory = tempfile::tempdir().unwrap();
2644        let layout = RepoLayout::single(directory.path());
2645        let url = "mkit+https://host/Uppercase";
2646        let error = super::open_with_config(url, &Config::default(), &layout)
2647            .err()
2648            .expect("uppercase repository identity must be rejected");
2649        assert!(
2650            matches!(error, super::DispatchError::MalformedUrl(_)),
2651            "expected malformed URL, got {error:?}"
2652        );
2653        let message = error.to_string();
2654        assert!(message.contains("Uppercase"), "{message}");
2655        assert!(!message.contains("invalid response"), "{message}");
2656    }
2657
2658    #[test]
2659    fn batch_lens_empty_input_is_no_batches() {
2660        assert_eq!(
2661            super::size_capped_batch_lens(&[], 1000, 10),
2662            Vec::<usize>::new()
2663        );
2664    }
2665
2666    #[test]
2667    fn batch_lens_fits_one_batch_under_both_caps() {
2668        assert_eq!(
2669            super::size_capped_batch_lens(&[10, 10, 10], 1000, 10),
2670            vec![3]
2671        );
2672    }
2673
2674    #[test]
2675    fn batch_lens_splits_on_byte_cap() {
2676        // Each item is 40 bytes; a 100-byte cap fits 2 per batch (80 <=
2677        // 100), not 3 (120 > 100) — entry cap (10) never binds here.
2678        assert_eq!(
2679            super::size_capped_batch_lens(&[40, 40, 40, 40, 40], 100, 10),
2680            vec![2, 2, 1]
2681        );
2682    }
2683
2684    #[test]
2685    fn batch_lens_splits_on_entry_cap() {
2686        // Every item is tiny (1 byte), so the byte cap (10_000) never
2687        // binds — only the entry cap (2) does.
2688        assert_eq!(
2689            super::size_capped_batch_lens(&[1, 1, 1, 1, 1], 10_000, 2),
2690            vec![2, 2, 1]
2691        );
2692    }
2693
2694    #[test]
2695    fn batch_lens_single_oversized_item_gets_its_own_batch() {
2696        // A lone item over `payload_cap` must not loop forever or be
2697        // dropped — mirrors `should_seal`'s "never true for an empty
2698        // writer" escape hatch.
2699        assert_eq!(super::size_capped_batch_lens(&[500], 100, 10), vec![1]);
2700        assert_eq!(
2701            super::size_capped_batch_lens(&[500, 10, 10], 100, 10),
2702            vec![1, 2]
2703        );
2704    }
2705
2706    // =================================================================
2707    // `prepare_raw_batch` / `prepare_delta_batch` — order preservation
2708    // and error propagation across BOTH the sequential and the rayon
2709    // branch. `pack_fanout_threshold()` is read directly (rather than
2710    // hardcoding a count) so a chunk built at exactly that length
2711    // reliably lands in the parallel branch regardless of how many
2712    // cores the test happens to run on.
2713    // =================================================================
2714
2715    fn store() -> (tempfile::TempDir, ObjectStore) {
2716        let dir = tempfile::tempdir().expect("tempdir");
2717        let layout = RepoLayout::single(dir.path());
2718        let store = ObjectStore::init(&layout).expect("init store");
2719        (dir, store)
2720    }
2721
2722    #[test]
2723    fn prepare_raw_batch_preserves_order_sequential_and_parallel() {
2724        let (_dir, store) = store();
2725        for &n in &[1usize, super::pack_fanout_threshold()] {
2726            let hashes: Vec<mkit_core::hash::Hash> = (0..n)
2727                .map(|i| store.write(format!("raw entry #{i}").as_bytes()).unwrap())
2728                .collect();
2729            let prepared = super::prepare_raw_batch(&store, &hashes).expect("prepare raw batch");
2730            assert_eq!(
2731                prepared.iter().map(PreparedRaw::hash).collect::<Vec<_>>(),
2732                hashes,
2733                "prepare_raw_batch must preserve input order at n={n}"
2734            );
2735        }
2736    }
2737
2738    #[test]
2739    fn prepare_raw_batch_propagates_a_missing_object_error_sequential_and_parallel() {
2740        let (_dir, store) = store();
2741        for &n in &[1usize, super::pack_fanout_threshold()] {
2742            let mut hashes: Vec<mkit_core::hash::Hash> = (0..n.saturating_sub(1))
2743                .map(|i| store.write(format!("raw entry #{i}").as_bytes()).unwrap())
2744                .collect();
2745            // A hash of content never written to the store — `store.read`
2746            // must fail with `ObjectNotFound`, and that failure must
2747            // surface through `prepare_raw_batch` rather than being
2748            // silently dropped or panicking on either branch.
2749            hashes.push(mkit_core::hash::hash(b"never written"));
2750            let err = super::prepare_raw_batch(&store, &hashes).expect_err(
2751                "a missing object must fail prepare_raw_batch instead of being dropped",
2752            );
2753            assert!(matches!(
2754                err,
2755                super::DispatchError::Store(mkit_core::store::StoreError::ObjectNotFound(_))
2756            ));
2757        }
2758    }
2759
2760    #[test]
2761    fn prepare_delta_batch_preserves_order_sequential_and_parallel() {
2762        use mkit_core::transfer::PlannedDelta;
2763        for &n in &[1usize, super::pack_fanout_threshold()] {
2764            let deltas: Vec<PlannedDelta> = (0..n)
2765                .map(|i| PlannedDelta {
2766                    target: mkit_core::hash::hash(format!("target #{i}").as_bytes()),
2767                    base: mkit_core::hash::hash(format!("base #{i}").as_bytes()),
2768                    stream: format!("delta stream #{i}").into_bytes(),
2769                })
2770                .collect();
2771            let expected_bases: Vec<_> = deltas.iter().map(|d| d.base).collect();
2772            let prepared = super::prepare_delta_batch(deltas);
2773            assert_eq!(
2774                prepared.iter().map(PreparedDelta::base).collect::<Vec<_>>(),
2775                expected_bases,
2776                "prepare_delta_batch must preserve input order at n={n}"
2777            );
2778        }
2779    }
2780
2781    #[test]
2782    fn encode_delta_candidates_batch_preserves_order_sequential_and_parallel() {
2783        let (_dir, store) = store();
2784        for &n in &[
2785            1usize,
2786            super::delta_plan_fanout_threshold(),
2787            super::DELTA_CANDIDATE_BATCH_CAP * 2 + 1,
2788        ] {
2789            // Each candidate's base/target pair is a small, localized edit
2790            // of distinct-per-index content — a real delta candidate (the
2791            // encoded stream is smaller than the raw target), and distinct
2792            // across candidates so no candidate's result could accidentally
2793            // match another's and mask an ordering bug.
2794            let candidates: Vec<transfer::DeltaCandidate> = (0..n)
2795                .map(|i| {
2796                    let mut base_bytes =
2797                        format!("encode-delta-candidates-batch fixture #{i}\n").into_bytes();
2798                    base_bytes.extend_from_slice(
2799                        b"the quick brown fox jumps over the lazy dog\n"
2800                            .repeat(4)
2801                            .as_slice(),
2802                    );
2803                    let base = store.write(&base_bytes).unwrap();
2804                    let mut target_bytes = base_bytes.clone();
2805                    target_bytes.extend_from_slice(b"-- edited --\n");
2806                    let target = store.write(&target_bytes).unwrap();
2807                    transfer::DeltaCandidate { target, base }
2808                })
2809                .collect();
2810
2811            // Ground truth: encode each candidate one at a time, in order.
2812            let expected: Vec<Option<Vec<u8>>> = candidates
2813                .iter()
2814                .map(|&c| {
2815                    transfer::encode_delta_candidate(&store, c)
2816                        .expect("encode delta candidate")
2817                        .map(|p| p.stream)
2818                })
2819                .collect();
2820            assert!(
2821                expected.iter().all(Option::is_some),
2822                "every fixture candidate must actually delta-encode smaller than raw"
2823            );
2824
2825            let actual: Vec<Option<Vec<u8>>> =
2826                super::encode_delta_candidates_batch(&store, &candidates)
2827                    .expect("encode delta candidates batch")
2828                    .into_iter()
2829                    .map(|r| r.map(|p| p.stream))
2830                    .collect();
2831            assert_eq!(
2832                actual, expected,
2833                "encode_delta_candidates_batch must preserve input order at n={n}"
2834            );
2835        }
2836    }
2837
2838    #[test]
2839    fn encode_delta_candidates_batch_propagates_a_missing_object_error_sequential_and_parallel() {
2840        let (_dir, store) = store();
2841        for &n in &[
2842            1usize,
2843            super::delta_plan_fanout_threshold(),
2844            super::DELTA_CANDIDATE_BATCH_CAP * 2 + 1,
2845        ] {
2846            let mut candidates: Vec<transfer::DeltaCandidate> = (0..n.saturating_sub(1))
2847                .map(|i| {
2848                    let base = store
2849                        .write(format!("present base #{i}").as_bytes())
2850                        .unwrap();
2851                    let target = store
2852                        .write(format!("present target #{i}").as_bytes())
2853                        .unwrap();
2854                    transfer::DeltaCandidate { target, base }
2855                })
2856                .collect();
2857            // A target hash never written to the store — `store.read` must
2858            // fail with `ObjectNotFound`, surfaced through
2859            // `encode_delta_candidates_batch` rather than dropped or
2860            // panicking, on either branch.
2861            candidates.push(transfer::DeltaCandidate {
2862                target: mkit_core::hash::hash(b"never written target"),
2863                base: mkit_core::hash::hash(b"never written base"),
2864            });
2865            let err = super::encode_delta_candidates_batch(&store, &candidates).expect_err(
2866                "a missing object must fail encode_delta_candidates_batch instead of being dropped",
2867            );
2868            assert!(matches!(
2869                err,
2870                mkit_core::store::StoreError::ObjectNotFound(_)
2871            ));
2872        }
2873    }
2874
2875    #[test]
2876    fn delta_base_cache_respects_bytes_and_skips_single_use_bases() {
2877        let (_dir, store) = store();
2878        let mut candidates = Vec::new();
2879        for (value, size, uses) in [
2880            (b'a', 128usize, 2usize),
2881            (b'b', 128, 2),
2882            (b'c', 1024, 2),
2883            (b'd', 64, 1),
2884        ] {
2885            let base_bytes = vec![value; size];
2886            let base = store.write(&base_bytes).unwrap();
2887            let mut target_bytes = base_bytes;
2888            target_bytes.extend_from_slice(b"edited");
2889            let target = store.write(&target_bytes).unwrap();
2890            candidates.extend(std::iter::repeat_n(
2891                transfer::DeltaCandidate { target, base },
2892                uses,
2893            ));
2894        }
2895        // Only one of the two repeated 128-byte bases fits. The 1 KiB
2896        // base exceeds the budget and the 64-byte base is used only once.
2897        let cache = super::cache_delta_bases(&store, &candidates, 192).unwrap();
2898        assert_eq!(cache.len(), 1);
2899        assert_eq!(cache.values().map(Vec::len).sum::<usize>(), 128);
2900    }
2901
2902    #[test]
2903    fn delta_base_cache_does_not_read_past_a_failed_batch() {
2904        let (_dir, store) = store();
2905        let base = store.write(b"present base").unwrap();
2906        let missing_target = mkit_core::hash::hash(b"missing first-batch target");
2907        let mut candidates = vec![
2908            transfer::DeltaCandidate {
2909                target: missing_target,
2910                base,
2911            };
2912            super::DELTA_CANDIDATE_BATCH_CAP
2913        ];
2914        // If all bases are prefetched for the whole push, this missing base
2915        // fails first. Bounded processing must encounter the first batch's
2916        // missing target before it even tries to read the later base.
2917        candidates.extend(
2918            [transfer::DeltaCandidate {
2919                target: missing_target,
2920                base: mkit_core::hash::hash(b"missing later-batch base"),
2921            }; 2],
2922        );
2923        let err = super::encode_delta_candidates_batch(&store, &candidates).unwrap_err();
2924        assert!(matches!(
2925            err,
2926            mkit_core::store::StoreError::ObjectNotFound(h) if h == mkit_core::hash::to_hex(&missing_target)
2927        ));
2928    }
2929
2930    /// The three `ssh.*` trust-pinning keys, when set in `Config`, must
2931    /// map 1:1 into the `SshOptions` carried to the spawned `ssh(1)`
2932    /// child. This is the producer half of issue #389: without it the
2933    /// keys are parsed but never reach the subprocess.
2934    #[test]
2935    fn populated_config_maps_to_ssh_options() {
2936        let cfg = Config {
2937            ssh_strict_host_key_checking: "yes".to_string(),
2938            ssh_user_known_hosts_file: "/path/to/project.known_hosts".to_string(),
2939            ssh_identity_file: "/path/to/id_ed25519".to_string(),
2940            ..Config::default()
2941        };
2942        let opts = ssh_options_from_config(&cfg);
2943        assert_eq!(opts.strict_host_key_checking, "yes");
2944        assert_eq!(opts.user_known_hosts_file, "/path/to/project.known_hosts");
2945        assert_eq!(opts.identity_file, "/path/to/id_ed25519");
2946    }
2947
2948    /// Empty `ssh.*` fields must map to empty `SshOptions` fields so
2949    /// `build_ssh_command` emits NO `-o`/`-i` flags and the user's
2950    /// `ssh(1)` defaults are inherited unchanged.
2951    #[test]
2952    fn empty_config_maps_to_empty_ssh_options() {
2953        let opts = ssh_options_from_config(&Config::default());
2954        assert!(opts.strict_host_key_checking.is_empty());
2955        assert!(opts.user_known_hosts_file.is_empty());
2956        assert!(opts.identity_file.is_empty());
2957    }
2958}