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}