pub(crate) mod applied_packs;
mod envelope_signer;
pub(crate) mod grants;
mod packmap;
mod split;
mod upload_receipts;
use mkit_core::layout::RepoLayout;
use std::path::Path;
use std::sync::Arc;
use applied_packs::AppliedPacks;
use mkit_core::hash::{HASH_LEN, Hash};
use mkit_core::object::Object;
use mkit_core::ops::merge::is_ancestor;
use mkit_core::ops::restore;
use mkit_core::pack::{self, PackError, PackWriter, PreparedDelta, PreparedRaw};
use mkit_core::protocol::{PackKey, Transport, TransportError, UploadLimits};
use mkit_core::refs::{self, Head};
use mkit_core::store::{ObjectStore, StoreError};
use mkit_core::transfer::{self, PackListError};
use mkit_transport_connect::admission::{AdmissionPolicy, is_reserved};
use mkit_transport_connect::{
ConnectTransport, PENDING_INTERRUPTED_MESSAGE, repository_identity_from_url,
};
use mkit_transport_file::FileTransport;
use mkit_transport_s3::S3Transport;
use mkit_transport_ssh::{SshInitError, SshOptions, SshTransport, parse_mkit_ssh_url};
use rayon::prelude::*;
pub use split::{MAX_CHAIN_COMMITS, MAX_SPLIT_STEPS, PushControl, StepAuthority, plan_push_steps};
use packmap::{
ChainAction, advance_packmap, apply_fetched_chain, commit_head, packmap_ref, probe_chain,
rebaseline_depth, resolve_and_download_chain,
};
const DEFAULT_REMOTE: &str = "default";
#[derive(Debug, thiserror::Error)]
pub enum DispatchError {
#[error("unsupported URL scheme: {0}")]
UnsupportedScheme(String),
#[error("malformed URL: {0}")]
MalformedUrl(String),
#[error(
"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)"
)]
RepositoryNotFound { identity: String, origin: String },
#[error("no HEAD branch to push")]
NoHead,
#[error("interrupted")]
Interrupted,
#[error("{0}")]
UploadInterrupted(String),
#[error("transport: {0}")]
Transport(#[from] TransportError),
#[error("refs: {0}")]
Refs(#[from] refs::RefError),
#[error("repo lock: {0}")]
RepoLock(#[from] mkit_core::repo_lock::LockError),
#[error("worktree discovery: {0}")]
Discover(#[from] mkit_core::layout::DiscoverError),
#[error("io: {0}")]
Io(#[from] std::io::Error),
#[error("store: {0}")]
Store(#[from] StoreError),
#[error("pack: {0}")]
Pack(#[from] PackError),
#[error("packlist: {0}")]
PackList(#[from] PackListError),
#[error("ssh init: {0}")]
SshInit(#[from] SshInitError),
#[error("pull requires HEAD to point at a branch")]
DetachedHead,
#[error("remote branch '{0}' not found")]
RemoteBranchMissing(String),
#[error("pull would not fast-forward branch '{branch}'; merge or rebase first")]
NonFastForwardPull { branch: String },
#[error("restore safety: {0}")]
RestoreSafety(String),
#[error("object is not a commit")]
NotCommit,
#[error("restore: {0}")]
Restore(#[from] restore::RestoreError),
#[error("{0}")]
UntrustedRemote(String),
#[error(
"updates were rejected for branch '{branch}' (non-fast-forward); fetch and merge first, or re-run with --force-with-lease / --force"
)]
NonFastForwardPush { branch: String },
#[error(
"could not establish the pack map for branch '{branch}' under concurrent pushes; retry"
)]
PackmapContended { branch: String },
#[error("pack map chain for branch '{branch}' is malformed (too deep, cyclic, or unreadable)")]
PackChainInvalid { branch: String },
#[error("remote advertised branch '{0}' but no pack map to reconstruct it")]
PackmapMissing(String),
#[error("remote advertised pack {pack} for branch '{branch}' but does not hold it")]
AdvertisedPackMissing { branch: String, pack: String },
#[error("remote is missing object {0} needed to reconstruct the ref")]
RemoteMissingObject(String),
#[error(
"fetched history is too large to verify (closure exceeds the {0}-object cap); refusing to publish an unverified ref"
)]
ClosureTooLarge(usize),
#[error("object {hash} failed signature verification: {reason}")]
UnsignedOrInvalidObject { hash: String, reason: String },
#[error("invalid remote name for applied-packs record: '{0}'")]
InvalidRemoteName(String),
#[error(
"push needs {packs} data packs but this server permits {limit} per advance{}{}; ask the operator to raise max_pack_bytes",
.commit.as_ref().map_or_else(String::new, |c| format!(" ({c} cannot be split further)")),
if *.holds_remote_head {
" (it is the first commit that contains the remote head: rebase onto the remote head, or merge it in a smaller or earlier commit)"
} else {
""
}
)]
PushTooLarge {
packs: usize,
limit: usize,
commit: Option<String>,
holds_remote_head: bool,
},
#[error(
"{reason}, and the retry needs {packs} data packs but this server permits {limit} per advance; nothing was uploaded for the retry"
)]
RetryTooLarge {
reason: RetryReason,
packs: usize,
limit: usize,
},
#[error("{0}")]
PushSplitLimit(String),
#[error("{0}")]
PushNotAuthorized(String),
#[error(
"{cause}; {published} of {total} advances were published and the last published advance was {head} on branch '{branch}'{}",
resume_hint(.cause)
)]
SplitInterrupted {
branch: String,
head: String,
published: usize,
total: usize,
cause: Box<DispatchError>,
},
#[error("upload ticket was rejected after the push was retried")]
TicketRejected,
#[error("delta base is unavailable after a restart")]
DeltaBaseUnavailable,
#[error("packlist still names a pack absent from this repository after retry")]
PacklistNotInRepository,
}
fn resume_hint(cause: &DispatchError) -> &'static str {
if matches!(cause, DispatchError::NonFastForwardPush { .. }) {
"; the branch was moved by another push, so fetch and merge before pushing again"
} else {
"; re-run the push to resume"
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
pub enum RetryReason {
#[error(
"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)"
)]
LostDeltaBase,
#[error(
"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)"
)]
Replanned,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PublishedPrefix {
pub branch: String,
pub published: usize,
pub total: usize,
pub head: String,
pub resumable: bool,
}
impl PublishedPrefix {
#[must_use]
pub fn note(&self) -> String {
format!(
"{} of {} advances were published; the last published advance on branch '{}' was {}{}",
self.published,
self.total,
self.branch,
self.head,
if self.resumable {
"; re-run the push to resume"
} else {
"; the branch was moved by another push, so fetch and merge before pushing again"
}
)
}
}
impl DispatchError {
#[must_use]
pub fn into_published_prefix(self) -> (Self, Option<PublishedPrefix>) {
match self {
Self::SplitInterrupted {
branch,
head,
published,
total,
cause,
} => {
let resumable = !matches!(*cause, Self::NonFastForwardPush { .. });
(
*cause,
Some(PublishedPrefix {
branch,
published,
total,
head,
resumable,
}),
)
}
other => (other, None),
}
}
}
fn repository_operation_error(tx: &dyn Transport, error: TransportError) -> DispatchError {
if let TransportError::RemoteError(message) = &error
&& message.starts_with("upload interrupted; ")
{
return DispatchError::UploadInterrupted(message.clone());
}
if matches!(&error, TransportError::RemoteError(message) if message == PENDING_INTERRUPTED_MESSAGE)
{
return DispatchError::Interrupted;
}
if matches!(&error, TransportError::PackNotFound)
&& let Some(address) = tx.repository_address()
{
return DispatchError::RepositoryNotFound {
identity: address.repository.to_owned(),
origin: address.origin.to_owned(),
};
}
error.into()
}
pub fn open_trusted(
endpoint: &str,
remote_name: &str,
repo_chosen: bool,
cfg: &crate::config::LayeredConfig,
layout: &RepoLayout,
) -> Result<Arc<dyn Transport>, DispatchError> {
crate::config::endpoint_credential_trust(cfg, endpoint, repo_chosen)
.map_err(DispatchError::UntrustedRemote)?;
open_with_config_for_remote(endpoint, &cfg.merged, layout, Some(remote_name))
}
pub(crate) struct PushRemote {
pub tx: Arc<dyn Transport>,
pub authority: Option<Arc<dyn StepAuthority>>,
}
pub(crate) fn open_trusted_for_push(
endpoint: &str,
remote_name: &str,
repo_chosen: bool,
cfg: &crate::config::LayeredConfig,
layout: &RepoLayout,
) -> Result<PushRemote, DispatchError> {
crate::config::endpoint_credential_trust(cfg, endpoint, repo_chosen)
.map_err(DispatchError::UntrustedRemote)?;
if !is_connect_url(endpoint) {
return Ok(PushRemote {
tx: open_with_config_for_remote(endpoint, &cfg.merged, layout, Some(remote_name))?,
authority: None,
});
}
let parts = open_connect_parts(endpoint, &cfg.merged, layout, Some(remote_name), true)?;
let tx = Arc::new(parts.transport);
let authority = parts
.signer_key
.zip(parts.grants)
.map(|(signer_key, grants)| {
Arc::new(grants::ConnectAuthority {
transport: tx.clone(),
grants,
signer_key,
}) as Arc<dyn StepAuthority>
});
Ok(PushRemote { tx, authority })
}
pub(crate) fn open_with_config(
url: &str,
cfg: &crate::config::Config,
layout: &RepoLayout,
) -> Result<Arc<dyn Transport>, DispatchError> {
open_with_config_for_remote(url, cfg, layout, None)
}
fn open_with_config_for_remote(
url: &str,
cfg: &crate::config::Config,
layout: &RepoLayout,
remote_name: Option<&str>,
) -> Result<Arc<dyn Transport>, DispatchError> {
if is_connect_url(url) {
return Ok(Arc::new(open_connect_with_config(
url,
cfg,
layout,
remote_name,
true,
)?));
}
open_with_ssh_options(url, &ssh_options_from_config(cfg), None)
}
fn is_connect_url(url: &str) -> bool {
url.starts_with("mkit+https://") || url.starts_with("mkit+http://")
}
pub(crate) fn open_connect_trusted(
endpoint: &str,
remote_name: &str,
repo_chosen: bool,
cfg: &crate::config::LayeredConfig,
layout: &RepoLayout,
sign: bool,
) -> Result<ConnectTransport, DispatchError> {
crate::config::endpoint_credential_trust(cfg, endpoint, repo_chosen)
.map_err(DispatchError::UntrustedRemote)?;
if !is_connect_url(endpoint) {
return Err(DispatchError::UnsupportedScheme(format!(
"`{endpoint}` is not an mkit+https:// or mkit+http:// remote; grant epochs and repository visibility are served over Connect"
)));
}
open_connect_with_config(endpoint, &cfg.merged, layout, Some(remote_name), sign)
}
pub(crate) fn open_connect_with_config(
url: &str,
cfg: &crate::config::Config,
layout: &RepoLayout,
remote_name: Option<&str>,
sign: bool,
) -> Result<ConnectTransport, DispatchError> {
Ok(open_connect_parts(url, cfg, layout, remote_name, sign)?.transport)
}
struct ConnectParts {
transport: ConnectTransport,
grants: Option<Arc<grants::LocalGrants>>,
signer_key: Option<String>,
}
fn open_connect_parts(
url: &str,
cfg: &crate::config::Config,
layout: &RepoLayout,
remote_name: Option<&str>,
sign: bool,
) -> Result<ConnectParts, DispatchError> {
let envelope_signer = if sign {
if cfg.transport_auth_envelope() && cfg.trusted_remote_endpoint.trim() != url {
return Err(DispatchError::UntrustedRemote(format!(
"refusing request signing for untrusted destination `{url}`; run `mkit config trusted_remote_endpoint {url}` before using ambient signing identity"
)));
}
envelope_signer_from_config(cfg, layout)?
} else {
None
};
validate_connect_repository(url)?;
let signed = envelope_signer.is_some();
let signer_key = envelope_signer
.as_ref()
.map(|signer| signer.public_key_hex());
let ca_file = if url.starts_with("mkit+https://")
&& std::env::var_os(mkit_transport_connect::tls::CA_FILE_ENV).is_none()
{
cfg.ssl_ca_file_path(layout).map_err(|error| {
DispatchError::Transport(TransportError::TlsConfiguration(error.to_string()))
})?
} else {
None
};
let mut tx = ConnectTransport::connect_with_signer_and_ca_file(
url,
envelope_signer,
ca_file.as_deref(),
)?
.with_receipt_store(Arc::new(upload_receipts::FilePartReceiptStore::new(
layout.upload_parts_dir(),
)))
.with_pending_observer(|event| {
crate::progress::pending_event(event);
!crate::signal::is_shutdown()
})
.with_upload_observer(|event| {
crate::progress::upload_event(event);
!crate::signal::is_shutdown()
})
.with_admission_receipt_observer(|receipt| {
eprintln!(
"note: remote returned a {} receipt for {}",
receipt.header,
receipt
.procedure
.rsplit('/')
.next()
.unwrap_or(receipt.procedure)
);
});
let grants = signed.then(|| Arc::new(grants::LocalGrants::from_stored(load_user_grants(cfg))));
if let Some(grants) = &grants {
tx = tx.with_grant_source(grants.clone());
}
if !cfg.admission_helper.is_empty() && cfg.trusted_remote_endpoint.trim() == url {
let responder = Arc::new(crate::admission_helper::ExecResponder {
path: cfg.admission_helper.clone().into(),
});
let mut policy = AdmissionPolicy::new(responder);
if let Some(name) = remote_name
&& let Some(headers) = cfg.remote_admission_headers.get(name)
{
let bearer = std::env::var("MKIT_API_TOKEN").is_ok_and(|s| !s.is_empty());
for header in headers.split(',').map(str::trim).filter(|s| !s.is_empty()) {
if is_reserved(header, bearer) {
eprintln!(
"warning: ignoring reserved header `{header}` in remote.{name}.admission_headers (see SPEC-TRANSPORT-CONNECT §5.1)"
);
} else {
policy = policy.with_extra_allowed(header);
}
}
}
tx = tx.with_admission(policy);
}
Ok(ConnectParts {
transport: tx,
grants,
signer_key,
})
}
fn load_user_grants(cfg: &crate::config::Config) -> Vec<crate::grants::store::StoredGrant> {
let rps = crate::grants::parse_relying_parties(&cfg.grant_webauthn_rp).unwrap_or_else(|e| {
eprintln!("warning: grant.webauthn_rp: {e}; treating no relying party as pinned");
Vec::new()
});
let store = match crate::grants::store::GrantStore::open_default() {
Ok(store) => store,
Err(e) => {
eprintln!("warning: {e}; using no stored grants");
return Vec::new();
}
};
let report = store.load(&rps);
for warning in &report.warnings {
eprintln!("warning: {warning}");
}
report.grants
}
pub(crate) fn owner_ed25519_signer(
cfg: &crate::config::Config,
layout: &RepoLayout,
) -> Result<Arc<dyn mkit_transport_connect::EnvelopeSigner>, String> {
match cfg.signer.as_str() {
"" | "legacy" => {
let key_path = crate::config::resolve_key_path(layout, &cfg.signing_key)
.map_err(|e| format!("signing_key: {e}"))?;
if !key_path.exists() {
return Err(format!(
"no signing key at {} — run `mkit keygen` first",
key_path.display()
));
}
let kp = mkit_core::sign::load_key(&key_path).map_err(|e| format!("load key: {e}"))?;
Ok(Arc::new(envelope_signer::RepoKeyEnvelopeSigner::new(kp)))
}
"keystore" => Ok(Arc::new(envelope_signer::KeystoreEnvelopeSigner::open(
cfg,
)?)),
other => Err(format!(
"unknown signer `{other}` — expected `legacy` or `keystore`"
)),
}
}
pub(crate) fn envelope_signer_from_config(
cfg: &crate::config::Config,
layout: &RepoLayout,
) -> Result<Option<Arc<dyn mkit_transport_connect::EnvelopeSigner>>, DispatchError> {
if !cfg.transport_auth_envelope() {
return Ok(None);
}
let remote_error = |msg: String| DispatchError::Transport(TransportError::RemoteError(msg));
match cfg.signer.as_str() {
"" | "legacy" => {
let key_path =
crate::config::resolve_key_path(layout, &cfg.signing_key).map_err(|e| {
remote_error(format!("transport_auth = envelope: signing_key: {e}"))
})?;
if !key_path.exists() {
return Err(remote_error(format!(
"transport_auth = envelope requires a signing key at {} — run `mkit keygen` first",
key_path.display()
)));
}
let kp = mkit_core::sign::load_key(&key_path)
.map_err(|e| remote_error(format!("transport_auth = envelope: load key: {e}")))?;
Ok(Some(
Arc::new(envelope_signer::RepoKeyEnvelopeSigner::new(kp))
as Arc<dyn mkit_transport_connect::EnvelopeSigner>,
))
}
"keystore" => {
let signer = envelope_signer::KeystoreEnvelopeSigner::open(cfg)
.map_err(|e| remote_error(format!("transport_auth = envelope: {e}")))?;
Ok(Some(
Arc::new(signer) as Arc<dyn mkit_transport_connect::EnvelopeSigner>
))
}
other => Err(remote_error(format!(
"transport_auth = envelope: unknown signer `{other}` — expected `legacy` or `keystore`"
))),
}
}
fn ssh_options_from_config(cfg: &crate::config::Config) -> SshOptions {
SshOptions {
strict_host_key_checking: cfg.ssh_strict_host_key_checking.clone(),
user_known_hosts_file: cfg.ssh_user_known_hosts_file.clone(),
identity_file: cfg.ssh_identity_file.clone(),
}
}
pub fn open(url: &str) -> Result<Arc<dyn Transport>, DispatchError> {
open_with_ssh_options(url, &SshOptions::default(), None)
}
fn open_with_ssh_options(
url: &str,
ssh_options: &SshOptions,
envelope_signer: Option<Arc<dyn mkit_transport_connect::EnvelopeSigner>>,
) -> Result<Arc<dyn Transport>, DispatchError> {
if url.starts_with("git+") {
return Err(DispatchError::UnsupportedScheme(format!(
"'{url}' is a git-bridge remote — native push/pull/fetch/clone do not \
speak git transports; use `mkit git export` / `mkit git import` / \
`mkit git pull` (feature git-bridge)"
)));
}
if let Some(rest) = url.strip_prefix("mkit+file://") {
let path = Path::new(rest);
return Ok(Arc::new(FileTransport::new(path)));
}
if url.starts_with("mkit+memory://") {
return Err(DispatchError::UnsupportedScheme(
"mkit+memory:// must be driven via in-process harness (see tests)".to_string(),
));
}
if url.starts_with("mkit+https://") || url.starts_with("mkit+http://") {
validate_connect_repository(url)?;
let tx = ConnectTransport::connect_with_signer(url, envelope_signer)?;
return Ok(Arc::new(tx));
}
if url.starts_with("mkit+s3://") {
let tx = S3Transport::connect(url)?;
return Ok(Arc::new(tx));
}
if url.starts_with("mkit+ssh://") {
let target = parse_mkit_ssh_url(url).map_err(SshInitError::from)?;
let tx = SshTransport::connect_with_options(&target, ssh_options)?;
return Ok(Arc::new(tx));
}
#[cfg(feature = "enc-transport")]
if url.starts_with("mkit+enc://") {
return open_enc(url);
}
Err(DispatchError::MalformedUrl(url.to_string()))
}
fn validate_connect_repository(url: &str) -> Result<(), DispatchError> {
repository_identity_from_url(url).map_err(|reason| {
let path = url
.split_once("://")
.and_then(|(_, rest)| rest.split_once('/'))
.map_or("", |(_, path)| path)
.split(['?', '#'])
.next()
.unwrap_or("")
.trim_matches('/');
DispatchError::MalformedUrl(format!("repository identity `{path}` in {url}: {reason}"))
})?;
Ok(())
}
#[cfg(feature = "enc-transport")]
const ENC_CLIENT_KEY_ENV: &str = "MKIT_ENC_CLIENT_KEY";
#[cfg(feature = "enc-transport")]
fn open_enc(url: &str) -> Result<Arc<dyn Transport>, DispatchError> {
use mkit_transport_enc::url::parse_enc_url;
let target = parse_enc_url(url).map_err(DispatchError::Transport)?;
let sk = load_or_ephemeral_client_key()?;
let tx = mkit_transport_enc::connect_tcp(&target.host, target.port, &target.server_pubkey, sk)
.map_err(|e| DispatchError::Transport(TransportError::RemoteError(e.to_string())))?;
Ok(Arc::new(tx))
}
#[cfg(feature = "enc-transport")]
fn load_or_ephemeral_client_key()
-> Result<commonware_cryptography::ed25519::PrivateKey, DispatchError> {
use commonware_codec::DecodeExt as _;
use commonware_cryptography::ed25519::PrivateKey;
use zeroize::Zeroizing;
let map_err = |e: String| DispatchError::Transport(TransportError::RemoteError(e));
if let Some(path) = std::env::var_os(ENC_CLIENT_KEY_ENV).filter(|s| !s.is_empty()) {
let seed = mkit_core::sign::load_raw_32(std::path::Path::new(&path))
.map_err(|e| map_err(format!("load {ENC_CLIENT_KEY_ENV}: {e}")))?;
return PrivateKey::decode(seed.as_ref())
.map_err(|e| map_err(format!("client key construction failed: {e}")));
}
let mut secret = Zeroizing::new([0u8; 32]);
getrandom::fill(secret.as_mut()).map_err(|e| map_err(e.to_string()))?;
PrivateKey::decode(secret.as_ref()).map_err(|e| map_err(e.to_string()))
}
pub fn push_all(cwd: &Path, tx: &dyn Transport) -> Result<usize, DispatchError> {
push_all_with(cwd, tx, None, false, None).map(|pushed| pushed.refs)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Pushed {
pub refs: usize,
pub steps: usize,
}
pub fn push_all_with(
cwd: &Path,
tx: &dyn Transport,
remote: Option<&str>,
force: bool,
authority: Option<&dyn StepAuthority>,
) -> Result<Pushed, DispatchError> {
let layout = mkit_core::layout::discover(cwd)?;
let store = crate::commands::open_store_configured(&layout)?;
let refs_list = crate::commands::list_refs_parallel(&layout)?;
let remote = remote.unwrap_or(DEFAULT_REMOTE);
let mut n = 0;
let mut steps = 0;
let shallow = shallow_boundaries(&layout)?;
let mut tracking = refs::RemoteRefBatch::new(&layout, remote)?;
let result: Result<(), DispatchError> = (|| {
for r in refs_list {
if crate::signal::is_shutdown() {
return Err(DispatchError::Interrupted);
}
let Some(h) = r.hash else { continue };
let condition = if force {
refs::RefWriteCondition::Any
} else {
match refs::read_remote_ref(&layout, remote, &r.name)? {
Some(tracked) => refs::RefWriteCondition::Match(tracked),
None => refs::RefWriteCondition::Missing,
}
};
let control = PushControl {
authority,
shallow: shallow.clone(),
..PushControl::default()
};
steps += push_branch_steps(
tx,
&store,
&r.name,
h,
condition,
rebaseline_depth(),
pack::MAX_TOTAL_PAYLOAD,
&control,
&mut |published| Ok(tracking.write(&r.name, &published)?),
)?;
n += 1;
}
Ok(())
})();
tracking.commit()?;
result?;
Ok(Pushed { refs: n, steps })
}
fn shallow_boundaries(
layout: &RepoLayout,
) -> Result<std::collections::HashSet<Hash>, DispatchError> {
Ok(refs::load_shallow_boundaries(layout)?
.unwrap_or_default()
.into_iter()
.collect())
}
pub fn is_fast_forward(cwd: &Path, old: Option<Hash>, new: Hash) -> Result<bool, DispatchError> {
match old {
None => Ok(true),
Some(o) if o == new => Ok(true),
Some(o) => {
let layout = mkit_core::layout::discover(cwd)?;
let store = crate::commands::open_store_configured(&layout)?;
Ok(is_ancestor(&store, o, new)?)
}
}
}
#[derive(Debug, Clone, Copy)]
pub enum PushLease {
Force,
WithLease,
FastForward,
}
pub fn lease_condition(
cwd: &Path,
remote: &str,
branch: &str,
lease: PushLease,
) -> Result<refs::RefWriteCondition, DispatchError> {
if matches!(lease, PushLease::Force) {
return Ok(refs::RefWriteCondition::Any);
}
let layout = mkit_core::layout::discover(cwd)?;
Ok(match refs::read_remote_ref(&layout, remote, branch)? {
Some(tracked) => refs::RefWriteCondition::Match(tracked),
None => refs::RefWriteCondition::Missing,
})
}
pub fn push_branch_tracked(
cwd: &Path,
tx: &dyn Transport,
remote: &str,
branch: &str,
remote_branch: &str,
lease: PushLease,
authority: Option<&dyn StepAuthority>,
) -> Result<(Hash, usize), DispatchError> {
let layout = mkit_core::layout::discover(cwd)?;
let store = crate::commands::open_store_configured(&layout)?;
let tip = refs::read_ref(&layout, branch)?
.ok_or_else(|| DispatchError::RemoteBranchMissing(branch.to_owned()))?;
if matches!(lease, PushLease::FastForward)
&& let Some(tracked) = refs::read_remote_ref(&layout, remote, remote_branch)?
&& !is_ancestor(&store, tracked, tip)?
{
return Err(DispatchError::NonFastForwardPush {
branch: remote_branch.to_owned(),
});
}
let condition = lease_condition(cwd, remote, remote_branch, lease)?;
let control = PushControl {
authority,
shallow: shallow_boundaries(&layout)?,
..PushControl::default()
};
let steps = push_branch_steps(
tx,
&store,
remote_branch,
tip,
condition,
rebaseline_depth(),
pack::MAX_TOTAL_PAYLOAD,
&control,
&mut |published| {
Ok(refs::write_remote_ref(
&layout,
remote,
remote_branch,
&published,
)?)
},
)?;
Ok((tip, steps))
}
pub fn push_branch(
tx: &dyn Transport,
store: &ObjectStore,
branch: &str,
tip: Hash,
condition: refs::RefWriteCondition,
) -> Result<(), DispatchError> {
push_branch_with_depth(tx, store, branch, tip, condition, rebaseline_depth())
}
pub fn push_branch_with_depth(
tx: &dyn Transport,
store: &ObjectStore,
branch: &str,
tip: Hash,
condition: refs::RefWriteCondition,
rebaseline_threshold: usize,
) -> Result<(), DispatchError> {
push_branch_with_limits(
tx,
store,
branch,
tip,
condition,
rebaseline_threshold,
pack::MAX_TOTAL_PAYLOAD,
)
}
pub fn push_branch_with_limits(
tx: &dyn Transport,
store: &ObjectStore,
branch: &str,
tip: Hash,
condition: refs::RefWriteCondition,
rebaseline_threshold: usize,
pack_payload_cap: u64,
) -> Result<(), DispatchError> {
push_branch_steps(
tx,
store,
branch,
tip,
condition,
rebaseline_threshold,
pack_payload_cap,
&PushControl::default(),
&mut |_| Ok(()),
)
.map(drop)
}
#[allow(clippy::too_many_arguments)]
pub fn push_branch_steps(
tx: &dyn Transport,
store: &ObjectStore,
branch: &str,
tip: Hash,
condition: refs::RefWriteCondition,
rebaseline_threshold: usize,
pack_payload_cap: u64,
control: &PushControl<'_>,
on_step: &mut dyn FnMut(Hash) -> Result<(), DispatchError>,
) -> Result<usize, DispatchError> {
let limits = tx.upload_limits();
let mut single = |plan| {
push_step_recovering(
tx,
store,
branch,
tip,
condition,
rebaseline_threshold,
pack_payload_cap,
plan,
)?;
on_step(tip)?;
Ok(1)
};
if limits.tickets_per_advance.is_none() {
return single(None);
}
refs::check_pushable_branch(branch)?;
let remote_tip = tx.read_ref(&format!("refs/heads/{branch}"))?;
let plan = transfer::plan_pack_with(store, tip, remote_tip, encode_delta_candidates_batch)?;
let cap = effective_payload_cap(pack_payload_cap, limits.max_pack_bytes)?;
if estimate_pack_sizes(store, &plan, cap, limits.max_pack_bytes)?.len()
<= split::data_pack_budget(limits)
{
return single(Some(plan));
}
let unknown_remote = remote_tip.is_some_and(|remote| !store.contains(&remote));
if unknown_remote {
return single(None);
}
match build_and_upload_packs(PackSink::Count, store, plan, cap, limits) {
Err(DispatchError::PushTooLarge { .. }) => {}
Err(DispatchError::Interrupted) => return Err(DispatchError::Interrupted),
_ => return single(None),
}
let steps = plan_push_steps(
store,
tip,
remote_tip,
limits,
pack_payload_cap,
control,
condition,
branch,
)?;
if steps.len() == 1 {
return single(None);
}
let total = steps.len();
let mut published: Option<Hash> = None;
for (index, step_tip) in steps.into_iter().enumerate() {
let step_condition = published.map_or(condition, refs::RefWriteCondition::Match);
let mut landed = false;
let outcome = (|| {
if crate::signal::is_shutdown() {
return Err(DispatchError::Interrupted);
}
crate::progress::report(crate::progress::Event::Step {
index: index + 1,
total,
});
push_step_recovering(
tx,
store,
branch,
step_tip,
step_condition,
rebaseline_threshold,
pack_payload_cap,
None,
)?;
landed = true;
on_step(step_tip)?;
crate::progress::step_committed(index + 1, total, &step_tip);
Ok(())
})();
if landed {
published = Some(step_tip);
}
if let Err(cause) = outcome {
return Err(match published {
Some(head) => DispatchError::SplitInterrupted {
branch: branch.to_owned(),
head: mkit_core::hash::to_hex(&head),
published: index + usize::from(landed),
total,
cause: Box::new(cause),
},
None => cause,
});
}
}
Ok(total)
}
#[allow(clippy::too_many_arguments)]
fn push_step_recovering(
tx: &dyn Transport,
store: &ObjectStore,
branch: &str,
tip: Hash,
condition: refs::RefWriteCondition,
rebaseline_threshold: usize,
pack_payload_cap: u64,
plan: Option<transfer::PackPlan>,
) -> Result<(), DispatchError> {
let once = |retry, plan| {
push_branch_once(
tx,
store,
branch,
tip,
condition,
rebaseline_threshold,
pack_payload_cap,
retry,
plan,
)
};
match once(None, plan) {
Err(DispatchError::TicketRejected | DispatchError::PacklistNotInRepository) => {
once(Some(RetryReason::Replanned), None)
}
Err(DispatchError::DeltaBaseUnavailable) => once(Some(RetryReason::LostDeltaBase), None),
result => result,
}
}
#[allow(clippy::too_many_arguments)]
fn push_branch_once(
tx: &dyn Transport,
store: &ObjectStore,
branch: &str,
tip: Hash,
condition: refs::RefWriteCondition,
rebaseline_threshold: usize,
pack_payload_cap: u64,
retry: Option<RetryReason>,
preplanned: Option<transfer::PackPlan>,
) -> Result<(), DispatchError> {
let force_self_contained = retry == Some(RetryReason::LostDeltaBase);
refs::check_pushable_branch(branch)?;
let mut plan = if let Some(plan) = preplanned {
plan
} else {
let remote_tip = tx.read_ref(&format!("refs/heads/{branch}"))?;
transfer::plan_pack_with(
store,
tip,
if force_self_contained {
None
} else {
remote_tip
},
encode_delta_candidates_batch,
)?
};
let limits = tx.upload_limits();
let effective_cap = effective_payload_cap(pack_payload_cap, limits.max_pack_bytes)?;
if plan.is_empty() {
return commit_head(tx, &format!("refs/heads/{branch}"), condition, &tip, branch);
}
if crate::signal::is_shutdown() {
return Err(DispatchError::Interrupted);
}
let mut rebaseline = false;
let mut resolved_chain = None;
if !force_self_contained
&& rebaseline_threshold > 0
&& let Some(pm) = tx.read_ref(&packmap_ref(branch))?
{
match probe_chain(tx, branch, pm) {
Ok(chain)
if chain.depth + 1 > rebaseline_threshold
&& tx.supports_atomic_advance()
&& !matches!(condition, refs::RefWriteCondition::Any) =>
{
rebaseline = true;
}
Ok(chain) => resolved_chain = Some(chain),
Err(DispatchError::PackChainInvalid { .. }) => {}
Err(e) => return Err(e),
}
}
if rebaseline {
let full = transfer::plan_pack_with(store, tip, None, encode_delta_candidates_batch)?;
let full_cap = effective_payload_cap(pack_payload_cap, limits.max_pack_bytes)?;
if estimated_ticketed_count(
&estimate_pack_sizes(store, &full, full_cap, limits.max_pack_bytes)?,
limits,
) > split::data_pack_budget(limits)
{
rebaseline = false;
} else {
plan = full;
}
}
if let Some(reason) = retry
&& estimated_ticketed_count(
&estimate_pack_sizes(store, &plan, effective_cap, limits.max_pack_bytes)?,
limits,
) > split::data_pack_budget(limits)
&& let Err(DispatchError::PushTooLarge { packs, limit, .. }) =
build_and_upload_packs(PackSink::Count, store, plan.clone(), effective_cap, limits)
{
return Err(DispatchError::RetryTooLarge {
reason,
packs,
limit,
});
}
let self_contained = plan.self_contained;
let head_ref = format!("refs/heads/{branch}");
let pack_keys = build_and_upload_packs(
PackSink::Upload {
tx,
head_ref: &head_ref,
},
store,
plan,
effective_cap,
limits,
)?;
let action = if rebaseline {
ChainAction::ResetSelfContained
} else {
ChainAction::Append { self_contained }
};
advance_packmap(
tx,
branch,
&pack_keys,
action,
resolved_chain,
condition,
tip,
)
}
fn effective_payload_cap(requested: u64, advertised: Option<u64>) -> Result<u64, DispatchError> {
let cap = requested
.min(pack::MAX_TOTAL_PAYLOAD)
.min(advertised.unwrap_or(pack::MAX_TOTAL_PAYLOAD));
if cap == 0 {
return Err(DispatchError::Transport(TransportError::PayloadTooLarge(0)));
}
Ok(cap)
}
fn serialized_bound(payload: u64, entries: usize) -> u64 {
payload
.saturating_add((entries as u64).saturating_mul(pack::ENTRY_FRAME_LEN as u64))
.saturating_add((pack::HEADER_LEN + pack::TRAILER_LEN) as u64)
}
fn ticketed_pack(limits: UploadLimits, serialized_size: u64) -> bool {
limits.tickets_per_advance.is_some()
&& serialized_size >= limits.ticket_threshold_bytes.unwrap_or(0)
}
fn estimated_ticketed_count(sizes: &[u64], limits: UploadLimits) -> usize {
sizes
.iter()
.filter(|&&size| ticketed_pack(limits, size))
.count()
}
fn estimate_pack_sizes(
store: &ObjectStore,
plan: &transfer::PackPlan,
payload_cap: u64,
max_pack_bytes: Option<u64>,
) -> Result<Vec<u64>, DispatchError> {
let mut packs = Vec::new();
let mut payload = 0_u64;
let mut entries = 0_usize;
let sizes = plan
.raw
.iter()
.map(|hash| store.object_metadata(hash).map(|meta| meta.len()))
.collect::<Result<Vec<_>, _>>()?;
for size in sizes.into_iter().chain(
plan.deltas
.iter()
.map(|delta| (HASH_LEN + delta.stream.len()) as u64),
) {
let next_payload = payload.saturating_add(size);
if entries > 0
&& (next_payload > payload_cap
|| max_pack_bytes
.is_some_and(|limit| serialized_bound(next_payload, entries + 1) > limit))
{
packs.push(serialized_bound(payload, entries));
payload = 0;
entries = 0;
}
payload = payload.saturating_add(size);
entries += 1;
}
if entries > 0 {
packs.push(serialized_bound(payload, entries));
}
Ok(packs)
}
#[derive(Clone, Copy)]
enum PackSink<'a> {
Upload {
tx: &'a dyn Transport,
head_ref: &'a str,
},
Count,
}
fn build_and_upload_packs(
sink: PackSink<'_>,
store: &ObjectStore,
plan: transfer::PackPlan,
payload_cap: u64,
limits: UploadLimits,
) -> Result<Vec<Hash>, DispatchError> {
let mut pack_keys = Vec::new();
let mut ticketed_count = 0;
let mut w = PackWriter::new();
let max_entries = pack_fanout_threshold().max(1);
let transfer::PackPlan {
raw, mut deltas, ..
} = plan;
let raw_sizes: Vec<u64> = raw
.iter()
.map(|h| Ok(store.object_metadata(h)?.len()))
.collect::<Result<_, DispatchError>>()?;
let mut start = 0;
for len in size_capped_batch_lens(&raw_sizes, payload_cap, max_entries) {
let chunk = &raw[start..start + len];
start += len;
for entry in prepare_raw_batch(store, chunk)? {
if should_seal(
&w,
entry.conservative_len() as u64,
payload_cap,
limits.max_pack_bytes,
) {
seal_pack(sink, &mut w, &mut pack_keys, &mut ticketed_count, limits)?;
}
w.push_prepared_raw(entry)?;
report_packed(sink);
}
}
let delta_sizes: Vec<u64> = deltas
.iter()
.map(|d| (HASH_LEN + d.stream.len()) as u64)
.collect();
for len in size_capped_batch_lens(&delta_sizes, payload_cap, max_entries) {
let chunk: Vec<transfer::PlannedDelta> = deltas.drain(..len).collect();
for entry in prepare_delta_batch(chunk) {
if should_seal(
&w,
entry.conservative_len() as u64,
payload_cap,
limits.max_pack_bytes,
) {
seal_pack(sink, &mut w, &mut pack_keys, &mut ticketed_count, limits)?;
}
w.push_prepared_delta(entry)?;
report_packed(sink);
}
}
seal_pack(sink, &mut w, &mut pack_keys, &mut ticketed_count, limits)?;
Ok(pack_keys)
}
fn report_packed(sink: PackSink<'_>) {
if matches!(sink, PackSink::Upload { .. }) {
crate::progress::report(crate::progress::Event::ObjectsPacked(1));
}
}
fn size_capped_batch_lens(sizes: &[u64], payload_cap: u64, max_entries: usize) -> Vec<usize> {
let mut lens = Vec::new();
let mut running: u64 = 0;
let mut count = 0usize;
for &size in sizes {
if count > 0 && (running.saturating_add(size) > payload_cap || count >= max_entries) {
lens.push(count);
running = 0;
count = 0;
}
running += size;
count += 1;
}
if count > 0 {
lens.push(count);
}
lens
}
const PACK_FANOUT_ENTRIES_PER_THREAD: usize = 8;
fn pack_fanout_threshold() -> usize {
crate::fanout::threshold(PACK_FANOUT_ENTRIES_PER_THREAD)
}
fn prepare_raw_batch(
store: &ObjectStore,
chunk: &[Hash],
) -> Result<Vec<PreparedRaw>, DispatchError> {
crate::fanout::try_map_seq_or_par(chunk, pack_fanout_threshold(), |h| {
Ok(PackWriter::prepare_raw(*h, store.read(h)?))
})
}
fn prepare_delta_batch(chunk: Vec<transfer::PlannedDelta>) -> Vec<PreparedDelta> {
if chunk.len() < pack_fanout_threshold() {
return chunk
.into_iter()
.map(|d| PackWriter::prepare_delta(d.base, d.stream))
.collect();
}
chunk
.into_par_iter()
.map(|d| PackWriter::prepare_delta(d.base, d.stream))
.collect()
}
const DELTA_PLAN_FANOUT_ENTRIES_PER_THREAD: usize = 1;
fn delta_plan_fanout_threshold() -> usize {
crate::fanout::threshold(DELTA_PLAN_FANOUT_ENTRIES_PER_THREAD)
}
fn encode_delta_candidates_batch(
store: &ObjectStore,
candidates: &[transfer::DeltaCandidate],
) -> Result<Vec<Option<transfer::PlannedDelta>>, StoreError> {
let mut results = Vec::with_capacity(candidates.len());
for batch in candidates.chunks(DELTA_CANDIDATE_BATCH_CAP) {
let base_bytes = cache_delta_bases(store, batch, DELTA_BASE_CACHE_BYTES)?;
results.extend(crate::fanout::try_map_seq_or_par(
batch,
delta_plan_fanout_threshold(),
|&c| match base_bytes.get(&c.base) {
Some(bytes) => transfer::encode_delta_candidate_with_base(store, c, bytes),
None => transfer::encode_delta_candidate(store, c),
},
)?);
}
Ok(results)
}
const DELTA_CANDIDATE_BATCH_CAP: usize = 64;
const DELTA_BASE_CACHE_BYTES: usize = 64 * 1024 * 1024;
fn cache_delta_bases(
store: &ObjectStore,
candidates: &[transfer::DeltaCandidate],
byte_budget: usize,
) -> Result<std::collections::HashMap<Hash, Vec<u8>>, StoreError> {
let mut counts = std::collections::BTreeMap::<Hash, usize>::new();
for candidate in candidates {
*counts.entry(candidate.base).or_default() += 1;
}
let mut cache = std::collections::HashMap::new();
let mut remaining = byte_budget;
for (base, count) in counts {
if count < 2 || remaining == 0 {
continue;
}
if store.object_metadata(&base)?.len() > remaining as u64 {
continue;
}
let bytes = store.read(&base)?;
if bytes.len() <= remaining {
remaining -= bytes.len();
cache.insert(base, bytes);
}
}
Ok(cache)
}
fn should_seal(
w: &PackWriter,
add_bound: u64,
payload_cap: u64,
max_pack_bytes: Option<u64>,
) -> bool {
let next_payload = w.total_payload().saturating_add(add_bound);
w.entry_count() > 0
&& (next_payload > payload_cap
|| max_pack_bytes
.is_some_and(|limit| serialized_bound(next_payload, w.entry_count() + 1) > limit))
}
fn seal_pack(
sink: PackSink<'_>,
w: &mut PackWriter,
pack_keys: &mut Vec<Hash>,
ticketed_count: &mut usize,
limits: UploadLimits,
) -> Result<(), DispatchError> {
if crate::signal::is_shutdown() {
return Err(DispatchError::Interrupted);
}
let sealed = std::mem::replace(w, PackWriter::new());
let pack = sealed.finish()?;
let is_ticketed = ticketed_pack(limits, pack.len() as u64);
let budget = split::data_pack_budget(limits);
if is_ticketed && *ticketed_count >= budget {
return Err(DispatchError::PushTooLarge {
packs: *ticketed_count + 1,
limit: budget,
commit: None,
holds_remote_head: false,
});
}
if limits
.max_pack_bytes
.is_some_and(|limit| pack.len() as u64 > limit)
{
return Err(DispatchError::Transport(TransportError::PayloadTooLarge(
pack.len(),
)));
}
let pack_key = pack::pack_key(&pack);
if let PackSink::Upload { tx, head_ref } = sink {
tx.upload_pack_via_ref(&pack, &PackKey::from_hash(pack_key), head_ref)
.map_err(|error| repository_operation_error(tx, error))?;
crate::progress::report(crate::progress::Event::PackUploaded(pack.len() as u64));
}
pack_keys.push(pack_key);
if is_ticketed {
*ticketed_count += 1;
}
Ok(())
}
pub fn pull_all(
cwd: &Path,
tx: &dyn Transport,
remote: &str,
target_branch: Option<&str>,
) -> Result<usize, DispatchError> {
pull_all_with(cwd, tx, remote, target_branch, true)
}
pub fn pull_all_with(
cwd: &Path,
tx: &dyn Transport,
remote: &str,
target_branch: Option<&str>,
require_signed: bool,
) -> Result<usize, DispatchError> {
let layout = mkit_core::layout::discover(cwd)?;
let store = crate::commands::open_store_configured(&layout)?;
let n = fetch_objects(&store, &layout, tx, remote, target_branch, require_signed)?;
let remote_refs = crate::commands::list_remote_refs_parallel(&layout, remote)?
.into_iter()
.filter_map(|r| r.hash.map(|hash| (r.name, hash)))
.collect::<Vec<_>>();
if remote_refs.is_empty() {
return Ok(n);
}
let _lock = mkit_core::repo_lock::acquire_default(
layout.worktree_state_dir(),
crate::commands::WORKTREE_LOCK,
)?;
crate::commands::warn_if_served(&layout);
let original_head = refs::read_head(&layout).ok();
let (branch, local_tip, remote_tip) = match &original_head {
Some(Head::Branch(head_branch)) => {
let want_branch = target_branch.unwrap_or(head_branch.as_str());
let local_tip = refs::read_ref(&layout, want_branch)?;
let selected = if local_tip.is_some() || target_branch.is_some() {
remote_refs
.iter()
.find(|(name, _)| name == want_branch)
.ok_or_else(|| DispatchError::RemoteBranchMissing(want_branch.to_owned()))?
} else {
remote_refs
.iter()
.find(|(name, _)| name == want_branch)
.unwrap_or(&remote_refs[0])
};
(selected.0.clone(), local_tip, selected.1)
}
Some(Head::Detached(_)) => return Err(DispatchError::DetachedHead),
None => (remote_refs[0].0.clone(), None, remote_refs[0].1),
};
let ref_condition = if let Some(local_tip) = local_tip {
if local_tip == remote_tip {
return Ok(n);
}
if !is_ancestor(&store, local_tip, remote_tip)? {
return Err(DispatchError::NonFastForwardPull { branch });
}
refs::RefWriteCondition::Match(local_tip)
} else {
refs::RefWriteCondition::Missing
};
let tree = load_tree_hash(&store, remote_tip)?;
crate::commands::ensure_restore_safe(&layout, &store, tree)
.map_err(DispatchError::RestoreSafety)?;
crate::commands::write_ref_recording_history(&layout, &branch, ref_condition, &remote_tip)?;
if let Err(e) = refs::write_head_branch(&layout, &branch) {
rollback_pull_ref(&layout, &branch, local_tip, remote_tip)?;
return Err(e.into());
}
if let Err(e) = crate::commands::restore_worktree_and_index(&layout, &store, tree) {
if let Err(rollback) =
rollback_pull_ref_and_head(&layout, &branch, local_tip, remote_tip, original_head)
{
return Err(DispatchError::RestoreSafety(format!(
"{e}; additionally failed to roll back ref: {rollback}"
)));
}
return Err(DispatchError::RestoreSafety(e));
}
Ok(n)
}
fn rollback_pull_ref_and_head(
layout: &RepoLayout,
branch: &str,
local_tip: Option<Hash>,
remote_tip: Hash,
original_head: Option<Head>,
) -> Result<(), String> {
rollback_pull_ref(layout, branch, local_tip, remote_tip).map_err(|e| e.to_string())?;
match original_head {
Some(Head::Branch(name)) => refs::write_head_branch(layout, &name),
Some(Head::Detached(hash)) => refs::write_head_detached(layout, &hash),
None => Ok(()),
}
.map_err(|e| e.to_string())
}
fn rollback_pull_ref(
layout: &RepoLayout,
branch: &str,
local_tip: Option<Hash>,
remote_tip: Hash,
) -> Result<(), refs::RefError> {
if let Some(local_tip) = local_tip {
crate::commands::write_ref_recording_history(
layout,
branch,
refs::RefWriteCondition::Match(remote_tip),
&local_tip,
)
} else if refs::read_ref(layout, branch)? == Some(remote_tip) {
refs::delete_ref(layout, branch)
} else {
Ok(())
}
}
pub fn fetch_all(cwd: &Path, tx: &dyn Transport, remote: &str) -> Result<usize, DispatchError> {
fetch_all_with(cwd, tx, remote, true)
}
pub fn fetch_all_with(
cwd: &Path,
tx: &dyn Transport,
remote: &str,
require_signed: bool,
) -> Result<usize, DispatchError> {
let layout = mkit_core::layout::discover(cwd)?;
let store = crate::commands::open_store_configured(&layout)?;
fetch_objects(&store, &layout, tx, remote, None, require_signed)
}
fn fetch_objects(
store: &ObjectStore,
layout: &RepoLayout,
tx: &dyn Transport,
remote: &str,
target_branch: Option<&str>,
require_signed: bool,
) -> Result<usize, DispatchError> {
let mut applied = AppliedPacks::load_or_empty(layout, remote);
let result = fetch_objects_inner(
store,
layout,
tx,
remote,
target_branch,
&mut applied,
require_signed,
);
persist_record(&mut applied, remote);
result
}
fn fetch_objects_inner(
store: &ObjectStore,
layout: &RepoLayout,
tx: &dyn Transport,
remote: &str,
target_branch: Option<&str>,
applied: &mut AppliedPacks,
require_signed: bool,
) -> Result<usize, DispatchError> {
let mut remote_refs = tx
.list_refs("refs/heads/")
.map_err(|error| repository_operation_error(tx, error))?;
if let Some(branch) = target_branch
&& !remote_refs.iter().any(|r| r.name == branch)
{
match tx.read_ref(&format!("refs/heads/{branch}"))? {
Some(hash) => remote_refs.push(refs::Ref {
name: branch.to_owned(),
hash: Some(hash),
}),
None => return Err(DispatchError::RemoteBranchMissing(branch.to_owned())),
}
}
let mut n = 0;
let mut tracking = refs::RemoteRefBatch::new(layout, remote)?;
let result: Result<(), DispatchError> = (|| {
for r in remote_refs {
if crate::signal::is_shutdown() {
return Err(DispatchError::Interrupted);
}
let Some(h) = r.hash else { continue };
let Some(chain_head) = tx.read_ref(&packmap_ref(&r.name))? else {
if tx.read_ref(&format!("refs/heads/{}", r.name))?.is_none() {
if crate::progress::should_report(false) {
eprintln!(
"fetch: skipping branch `{}`: listed by an eventual ListRefs but \
no longer present (deleted since the listing)",
r.name
);
}
continue;
}
return Err(DispatchError::PackmapMissing(r.name.clone()));
};
let fetched = resolve_and_download_chain(tx, &r.name, chain_head, applied)?;
let lock = mkit_core::repo_lock::acquire_default(
layout.worktree_state_dir(),
crate::commands::WORKTREE_LOCK,
)?;
crate::commands::warn_if_served(layout);
let (published_tip, _lock) = match apply_fetched_chain(
store,
tx,
remote,
&r.name,
fetched,
h,
applied,
require_signed,
) {
Ok(()) => (h, lock),
Err(e @ DispatchError::RemoteMissingObject(_)) => {
let (Some(fresh_h), Some(fresh_head)) = (
tx.read_ref(&format!("refs/heads/{}", r.name))?,
tx.read_ref(&packmap_ref(&r.name))?,
) else {
return Err(e);
};
drop(lock);
let fresh_fetched =
resolve_and_download_chain(tx, &r.name, fresh_head, applied)?;
let lock = mkit_core::repo_lock::acquire_default(
layout.worktree_state_dir(),
crate::commands::WORKTREE_LOCK,
)?;
crate::commands::warn_if_served(layout);
apply_fetched_chain(
store,
tx,
remote,
&r.name,
fresh_fetched,
fresh_h,
applied,
require_signed,
)?;
(fresh_h, lock)
}
Err(e) => return Err(e),
};
tracking.write(&r.name, &published_tip)?;
n += 1;
}
Ok(())
})();
tracking.commit()?;
result?;
Ok(n)
}
fn persist_record(applied: &mut AppliedPacks, remote: &str) {
if let Err(e) = applied.persist() {
eprintln!(
"warning: could not persist applied-packs record for remote '{remote}' ({e}); it will be rebuilt on the next fetch"
);
}
}
pub(crate) fn verify_closure_present(store: &ObjectStore, tip: &Hash) -> Result<(), DispatchError> {
match mkit_core::ops::reachable_closure_checked(store, std::iter::once(tip)) {
Ok((_, false)) => Ok(()),
Ok((_, true)) => Err(DispatchError::ClosureTooLarge(
mkit_core::ops::graph::MAX_REACHABLE,
)),
Err(StoreError::ObjectNotFound(hex)) => Err(DispatchError::RemoteMissingObject(hex)),
Err(e) => Err(e.into()),
}
}
fn load_tree_hash(store: &ObjectStore, commit_hash: Hash) -> Result<Hash, DispatchError> {
match store.read_object(&commit_hash)? {
Object::Commit(c) => Ok(c.tree_hash),
Object::Remix(r) => Ok(r.tree_hash),
_ => Err(DispatchError::NotCommit),
}
}
#[cfg(test)]
mod tests {
use super::ssh_options_from_config;
use crate::config::Config;
use mkit_core::layout::RepoLayout;
use mkit_core::pack::{PreparedDelta, PreparedRaw};
use mkit_core::store::ObjectStore;
use mkit_core::transfer;
#[test]
fn envelope_destination_trust_precedes_missing_key_resolution() {
let directory = tempfile::tempdir().unwrap();
let layout = RepoLayout::single(directory.path());
for endpoint in ["mkit+https://untrusted.example", "mkit+http://127.0.0.1:1"] {
let mut cfg = Config {
transport_auth: "envelope".into(),
signer: "legacy".into(),
signing_key: ".mkit/keys/missing-signing-key".into(),
..Config::default()
};
let error = super::open_with_config(endpoint, &cfg, &layout)
.err()
.expect("untrusted signing must fail");
assert!(
matches!(error, super::DispatchError::UntrustedRemote(_)),
"destination rejection must precede key access: {error}"
);
cfg.trusted_remote_endpoint = endpoint.into();
let error = super::open_with_config(endpoint, &cfg, &layout)
.err()
.expect("trusted endpoint now resolves the absent key");
assert!(
error.to_string().contains("requires a signing key"),
"trusted destination should reach key resolution: {error}"
);
}
}
#[test]
fn configured_connect_open_reports_malformed_repository_identity() {
let directory = tempfile::tempdir().unwrap();
let layout = RepoLayout::single(directory.path());
let url = "mkit+https://host/Uppercase";
let error = super::open_with_config(url, &Config::default(), &layout)
.err()
.expect("uppercase repository identity must be rejected");
assert!(
matches!(error, super::DispatchError::MalformedUrl(_)),
"expected malformed URL, got {error:?}"
);
let message = error.to_string();
assert!(message.contains("Uppercase"), "{message}");
assert!(!message.contains("invalid response"), "{message}");
}
#[test]
fn batch_lens_empty_input_is_no_batches() {
assert_eq!(
super::size_capped_batch_lens(&[], 1000, 10),
Vec::<usize>::new()
);
}
#[test]
fn batch_lens_fits_one_batch_under_both_caps() {
assert_eq!(
super::size_capped_batch_lens(&[10, 10, 10], 1000, 10),
vec![3]
);
}
#[test]
fn batch_lens_splits_on_byte_cap() {
assert_eq!(
super::size_capped_batch_lens(&[40, 40, 40, 40, 40], 100, 10),
vec![2, 2, 1]
);
}
#[test]
fn batch_lens_splits_on_entry_cap() {
assert_eq!(
super::size_capped_batch_lens(&[1, 1, 1, 1, 1], 10_000, 2),
vec![2, 2, 1]
);
}
#[test]
fn batch_lens_single_oversized_item_gets_its_own_batch() {
assert_eq!(super::size_capped_batch_lens(&[500], 100, 10), vec![1]);
assert_eq!(
super::size_capped_batch_lens(&[500, 10, 10], 100, 10),
vec![1, 2]
);
}
fn store() -> (tempfile::TempDir, ObjectStore) {
let dir = tempfile::tempdir().expect("tempdir");
let layout = RepoLayout::single(dir.path());
let store = ObjectStore::init(&layout).expect("init store");
(dir, store)
}
#[test]
fn prepare_raw_batch_preserves_order_sequential_and_parallel() {
let (_dir, store) = store();
for &n in &[1usize, super::pack_fanout_threshold()] {
let hashes: Vec<mkit_core::hash::Hash> = (0..n)
.map(|i| store.write(format!("raw entry #{i}").as_bytes()).unwrap())
.collect();
let prepared = super::prepare_raw_batch(&store, &hashes).expect("prepare raw batch");
assert_eq!(
prepared.iter().map(PreparedRaw::hash).collect::<Vec<_>>(),
hashes,
"prepare_raw_batch must preserve input order at n={n}"
);
}
}
#[test]
fn prepare_raw_batch_propagates_a_missing_object_error_sequential_and_parallel() {
let (_dir, store) = store();
for &n in &[1usize, super::pack_fanout_threshold()] {
let mut hashes: Vec<mkit_core::hash::Hash> = (0..n.saturating_sub(1))
.map(|i| store.write(format!("raw entry #{i}").as_bytes()).unwrap())
.collect();
hashes.push(mkit_core::hash::hash(b"never written"));
let err = super::prepare_raw_batch(&store, &hashes).expect_err(
"a missing object must fail prepare_raw_batch instead of being dropped",
);
assert!(matches!(
err,
super::DispatchError::Store(mkit_core::store::StoreError::ObjectNotFound(_))
));
}
}
#[test]
fn prepare_delta_batch_preserves_order_sequential_and_parallel() {
use mkit_core::transfer::PlannedDelta;
for &n in &[1usize, super::pack_fanout_threshold()] {
let deltas: Vec<PlannedDelta> = (0..n)
.map(|i| PlannedDelta {
target: mkit_core::hash::hash(format!("target #{i}").as_bytes()),
base: mkit_core::hash::hash(format!("base #{i}").as_bytes()),
stream: format!("delta stream #{i}").into_bytes(),
})
.collect();
let expected_bases: Vec<_> = deltas.iter().map(|d| d.base).collect();
let prepared = super::prepare_delta_batch(deltas);
assert_eq!(
prepared.iter().map(PreparedDelta::base).collect::<Vec<_>>(),
expected_bases,
"prepare_delta_batch must preserve input order at n={n}"
);
}
}
#[test]
fn encode_delta_candidates_batch_preserves_order_sequential_and_parallel() {
let (_dir, store) = store();
for &n in &[
1usize,
super::delta_plan_fanout_threshold(),
super::DELTA_CANDIDATE_BATCH_CAP * 2 + 1,
] {
let candidates: Vec<transfer::DeltaCandidate> = (0..n)
.map(|i| {
let mut base_bytes =
format!("encode-delta-candidates-batch fixture #{i}\n").into_bytes();
base_bytes.extend_from_slice(
b"the quick brown fox jumps over the lazy dog\n"
.repeat(4)
.as_slice(),
);
let base = store.write(&base_bytes).unwrap();
let mut target_bytes = base_bytes.clone();
target_bytes.extend_from_slice(b"-- edited --\n");
let target = store.write(&target_bytes).unwrap();
transfer::DeltaCandidate { target, base }
})
.collect();
let expected: Vec<Option<Vec<u8>>> = candidates
.iter()
.map(|&c| {
transfer::encode_delta_candidate(&store, c)
.expect("encode delta candidate")
.map(|p| p.stream)
})
.collect();
assert!(
expected.iter().all(Option::is_some),
"every fixture candidate must actually delta-encode smaller than raw"
);
let actual: Vec<Option<Vec<u8>>> =
super::encode_delta_candidates_batch(&store, &candidates)
.expect("encode delta candidates batch")
.into_iter()
.map(|r| r.map(|p| p.stream))
.collect();
assert_eq!(
actual, expected,
"encode_delta_candidates_batch must preserve input order at n={n}"
);
}
}
#[test]
fn encode_delta_candidates_batch_propagates_a_missing_object_error_sequential_and_parallel() {
let (_dir, store) = store();
for &n in &[
1usize,
super::delta_plan_fanout_threshold(),
super::DELTA_CANDIDATE_BATCH_CAP * 2 + 1,
] {
let mut candidates: Vec<transfer::DeltaCandidate> = (0..n.saturating_sub(1))
.map(|i| {
let base = store
.write(format!("present base #{i}").as_bytes())
.unwrap();
let target = store
.write(format!("present target #{i}").as_bytes())
.unwrap();
transfer::DeltaCandidate { target, base }
})
.collect();
candidates.push(transfer::DeltaCandidate {
target: mkit_core::hash::hash(b"never written target"),
base: mkit_core::hash::hash(b"never written base"),
});
let err = super::encode_delta_candidates_batch(&store, &candidates).expect_err(
"a missing object must fail encode_delta_candidates_batch instead of being dropped",
);
assert!(matches!(
err,
mkit_core::store::StoreError::ObjectNotFound(_)
));
}
}
#[test]
fn delta_base_cache_respects_bytes_and_skips_single_use_bases() {
let (_dir, store) = store();
let mut candidates = Vec::new();
for (value, size, uses) in [
(b'a', 128usize, 2usize),
(b'b', 128, 2),
(b'c', 1024, 2),
(b'd', 64, 1),
] {
let base_bytes = vec![value; size];
let base = store.write(&base_bytes).unwrap();
let mut target_bytes = base_bytes;
target_bytes.extend_from_slice(b"edited");
let target = store.write(&target_bytes).unwrap();
candidates.extend(std::iter::repeat_n(
transfer::DeltaCandidate { target, base },
uses,
));
}
let cache = super::cache_delta_bases(&store, &candidates, 192).unwrap();
assert_eq!(cache.len(), 1);
assert_eq!(cache.values().map(Vec::len).sum::<usize>(), 128);
}
#[test]
fn delta_base_cache_does_not_read_past_a_failed_batch() {
let (_dir, store) = store();
let base = store.write(b"present base").unwrap();
let missing_target = mkit_core::hash::hash(b"missing first-batch target");
let mut candidates = vec![
transfer::DeltaCandidate {
target: missing_target,
base,
};
super::DELTA_CANDIDATE_BATCH_CAP
];
candidates.extend(
[transfer::DeltaCandidate {
target: missing_target,
base: mkit_core::hash::hash(b"missing later-batch base"),
}; 2],
);
let err = super::encode_delta_candidates_batch(&store, &candidates).unwrap_err();
assert!(matches!(
err,
mkit_core::store::StoreError::ObjectNotFound(h) if h == mkit_core::hash::to_hex(&missing_target)
));
}
#[test]
fn populated_config_maps_to_ssh_options() {
let cfg = Config {
ssh_strict_host_key_checking: "yes".to_string(),
ssh_user_known_hosts_file: "/path/to/project.known_hosts".to_string(),
ssh_identity_file: "/path/to/id_ed25519".to_string(),
..Config::default()
};
let opts = ssh_options_from_config(&cfg);
assert_eq!(opts.strict_host_key_checking, "yes");
assert_eq!(opts.user_known_hosts_file, "/path/to/project.known_hosts");
assert_eq!(opts.identity_file, "/path/to/id_ed25519");
}
#[test]
fn empty_config_maps_to_empty_ssh_options() {
let opts = ssh_options_from_config(&Config::default());
assert!(opts.strict_host_key_checking.is_empty());
assert!(opts.user_known_hosts_file.is_empty());
assert!(opts.identity_file.is_empty());
}
}