use mkit_core::hash::{self, Hash};
use mkit_core::object::Object;
use mkit_core::pack::{self, PackReader};
use mkit_core::protocol::{AdvanceOutcome, PackKey, Transport, TransportError};
use mkit_core::refs;
use mkit_core::sign::{verify_commit, verify_remix, verify_tag};
use mkit_core::store::ObjectStore;
use mkit_core::transfer;
use super::DispatchError;
use super::applied_packs::AppliedPacks;
pub(crate) fn packmap_ref(branch: &str) -> String {
format!("refs/mkit/packmap/{branch}")
}
const PACKMAP_CAS_ATTEMPTS: u32 = 8;
const MAX_PACK_CHAIN_DEPTH: usize = 100_000;
const DEFAULT_REBASELINE_DEPTH: usize = 64;
pub(crate) fn rebaseline_depth() -> usize {
match std::env::var("MKIT_PACK_REBASELINE_DEPTH") {
Err(_) => DEFAULT_REBASELINE_DEPTH,
Ok(s) => s.parse::<usize>().unwrap_or_else(|_| {
eprintln!(
"warning: MKIT_PACK_REBASELINE_DEPTH='{s}' is not a valid non-negative \
integer; using the default {DEFAULT_REBASELINE_DEPTH} (set it to 0 to \
disable re-baselining)"
);
DEFAULT_REBASELINE_DEPTH
}),
}
}
fn download_packlist_node(
tx: &dyn Transport,
key: Hash,
) -> Result<transfer::PackListNode, DispatchError> {
let bytes = tx.download_blob(&PackKey::from_hash(key))?;
Ok(transfer::decode_packlist(&bytes)?)
}
fn walk_pack_chain(
tx: &dyn Transport,
branch: &str,
head_key: Hash,
) -> Result<Vec<transfer::PackListNode>, DispatchError> {
let invalid = || DispatchError::PackChainInvalid {
branch: branch.to_owned(),
};
let mut nodes = Vec::new();
let mut seen = std::collections::HashSet::new();
let mut cursor = Some(head_key);
while let Some(key) = cursor {
if crate::signal::is_shutdown() {
return Err(DispatchError::Interrupted);
}
if !seen.insert(key) || seen.len() > MAX_PACK_CHAIN_DEPTH {
return Err(invalid());
}
let node = match download_packlist_node(tx, key) {
Ok(n) => n,
Err(
DispatchError::Transport(TransportError::PackNotFound) | DispatchError::PackList(_),
) => return Err(invalid()),
Err(e) => return Err(e),
};
cursor = node.prev;
nodes.push(node);
}
Ok(nodes)
}
pub(crate) fn resolve_pack_chain(
tx: &dyn Transport,
branch: &str,
head_key: Hash,
) -> Result<Vec<Hash>, DispatchError> {
Ok(probe_chain(tx, branch, head_key)?.packs)
}
#[derive(Debug)]
pub(crate) struct ResolvedChain {
pub(crate) head: Hash,
pub(crate) depth: usize,
pub(crate) packs: Vec<Hash>,
}
pub(crate) fn probe_chain(
tx: &dyn Transport,
branch: &str,
head_key: Hash,
) -> Result<ResolvedChain, DispatchError> {
let mut nodes = walk_pack_chain(tx, branch, head_key)?;
let depth = nodes.len();
nodes.reverse(); let packs = nodes.into_iter().flat_map(|n| n.packs).collect();
Ok(ResolvedChain {
head: head_key,
depth,
packs,
})
}
#[derive(Debug, Clone, Copy)]
pub(crate) enum ChainAction {
Append {
self_contained: bool,
},
ResetSelfContained,
}
pub(crate) fn advance_packmap(
tx: &dyn Transport,
branch: &str,
pack_keys: &[Hash],
action: ChainAction,
resolved: Option<ResolvedChain>,
head_condition: refs::RefWriteCondition,
tip: Hash,
) -> Result<(), DispatchError> {
debug_assert!(
!pack_keys.is_empty(),
"advance_packmap requires at least one pack key — callers only reach this with a \
non-empty plan"
);
let packmap_name = packmap_ref(branch);
let head_name = format!("refs/heads/{branch}");
let mut cached = resolved;
for _ in 0..PACKMAP_CAS_ATTEMPTS {
if crate::signal::is_shutdown() {
return Err(DispatchError::Interrupted);
}
let prior = tx.read_ref(&packmap_name)?;
let prev = match action {
ChainAction::ResetSelfContained => {
debug_assert!(
tx.supports_atomic_advance(),
"re-baseline reset requires a transactional advance_refs (mkit #521)"
);
debug_assert!(
!matches!(head_condition, refs::RefWriteCondition::Any),
"re-baseline reset must not run with an `Any` head condition — the \
ordered advance_refs fallback would strand the head (mkit #521)"
);
None
}
ChainAction::Append { self_contained } => match prior {
None => None,
Some(p) => {
let packs = match cached.take() {
Some(c) if c.head == p => Ok(c.packs),
_ => resolve_pack_chain(tx, branch, p),
};
match packs {
Ok(packs) => {
let have: std::collections::HashSet<&Hash> = packs.iter().collect();
if pack_keys.iter().all(|k| have.contains(k)) {
return commit_head(tx, &head_name, head_condition, &tip, branch);
}
Some(p) }
Err(DispatchError::PackChainInvalid { .. }) if self_contained => None,
Err(e) => return Err(e),
}
}
},
};
let node = transfer::encode_packlist(prev, pack_keys)?;
let node_key = pack::pack_key(&node);
tx.upload_blob(&node, &PackKey::from_hash(node_key))?;
let packmap_condition = match prior {
Some(k) => refs::RefWriteCondition::Match(k),
None => refs::RefWriteCondition::Missing,
};
match tx.advance_refs(
&head_name,
head_condition,
&tip,
&packmap_name,
packmap_condition,
&node_key,
)? {
AdvanceOutcome::Committed => return Ok(()),
AdvanceOutcome::PackmapConflict => {}
AdvanceOutcome::HeadConflict => {
return Err(DispatchError::NonFastForwardPush {
branch: branch.to_owned(),
});
}
}
}
Err(DispatchError::PackmapContended {
branch: branch.to_owned(),
})
}
pub(crate) fn commit_head(
tx: &dyn Transport,
head_name: &str,
condition: refs::RefWriteCondition,
tip: &Hash,
branch: &str,
) -> Result<(), DispatchError> {
match tx.update_ref(head_name, condition, tip) {
Ok(()) => Ok(()),
Err(TransportError::RefConflict) => Err(DispatchError::NonFastForwardPush {
branch: branch.to_owned(),
}),
Err(e) => Err(e.into()),
}
}
pub(crate) struct FetchedChain {
chain: Vec<Hash>,
downloaded: Vec<(PackKey, Vec<u8>)>,
}
pub(crate) fn resolve_and_download_chain(
tx: &dyn Transport,
branch: &str,
head_key: Hash,
applied: &AppliedPacks,
) -> Result<FetchedChain, DispatchError> {
let chain = resolve_pack_chain(tx, branch, head_key)?;
let downloaded = download_pack_chain(tx, branch, &chain, applied)?;
Ok(FetchedChain { chain, downloaded })
}
pub(crate) fn apply_fetched_chain(
store: &ObjectStore,
tx: &dyn Transport,
remote: &str,
branch: &str,
fetched: FetchedChain,
tip: Hash,
applied: &mut AppliedPacks,
require_signed: bool,
) -> Result<(), DispatchError> {
let FetchedChain { chain, downloaded } = fetched;
let skipped = chain.len() - downloaded.len();
let stored = unpack_downloaded_packs(store, downloaded, applied)?;
match super::verify_closure_present(store, &tip) {
Ok(()) => verify_new_object_signatures(store, &stored, require_signed),
Err(e @ DispatchError::RemoteMissingObject(_)) if skipped > 0 => {
eprintln!(
"note: applied-packs record for remote '{remote}' branch '{branch}' looks stale ({e}); clearing it and re-fetching the full pack chain"
);
applied.clear();
let downloaded = download_pack_chain(tx, branch, &chain, applied)?;
let stored = unpack_downloaded_packs(store, downloaded, applied)?;
super::verify_closure_present(store, &tip)?;
verify_new_object_signatures(store, &stored, require_signed)
}
Err(e) => Err(e),
}
}
fn verify_new_object_signatures(
store: &ObjectStore,
stored: &[Hash],
require_signed: bool,
) -> Result<(), DispatchError> {
if !require_signed {
return Ok(());
}
for h in stored {
let obj = store.read_object(h)?;
let result = match &obj {
Object::Commit(c) => verify_commit(c),
Object::Remix(r) => verify_remix(r),
Object::Tag(t) => verify_tag(t),
Object::Blob(_) | Object::Tree(_) | Object::ChunkedBlob(_) | Object::Delta(_) => {
continue;
}
};
if let Err(e) = result {
return Err(DispatchError::UnsignedOrInvalidObject {
hash: hash::to_hex(h),
reason: e.to_string(),
});
}
}
Ok(())
}
fn download_pack_chain(
tx: &dyn Transport,
branch: &str,
chain: &[Hash],
applied: &AppliedPacks,
) -> Result<Vec<(PackKey, Vec<u8>)>, DispatchError> {
let mut out = Vec::new();
for &pk in chain {
if crate::signal::is_shutdown() {
return Err(DispatchError::Interrupted);
}
let key = PackKey::from_hash(pk);
if applied.contains(&key) {
continue;
}
let pack = match tx.download_pack(&key) {
Ok(b) => b,
Err(TransportError::PackNotFound) => {
return Err(DispatchError::AdvertisedPackMissing {
branch: branch.to_owned(),
pack: mkit_core::hash::to_hex(&pk),
});
}
Err(e) => return Err(e.into()),
};
out.push((key, pack));
}
Ok(out)
}
fn unpack_downloaded_packs(
store: &ObjectStore,
downloaded: Vec<(PackKey, Vec<u8>)>,
applied: &mut AppliedPacks,
) -> Result<Vec<Hash>, DispatchError> {
let mut stored = Vec::new();
for (key, pack) in downloaded {
if crate::signal::is_shutdown() {
return Err(DispatchError::Interrupted);
}
let report = PackReader::read(&pack, store)?;
let unpacked = (report.raw_count + report.delta_count) as usize;
if unpacked > 0 {
crate::progress::report(crate::progress::Event::ObjectsUnpacked(unpacked));
}
stored.extend(report.stored);
applied.insert(&key);
}
Ok(stored)
}
#[cfg(test)]
mod tests {
use super::*;
use mkit_core::hash;
use mkit_transport_memory::MemoryTransport;
fn h(seed: &str) -> Hash {
hash::hash(seed.as_bytes())
}
fn put_node(tx: &MemoryTransport, key: Hash, prev: Option<Hash>, packs: &[Hash]) {
let bytes = transfer::encode_packlist(prev, packs).unwrap();
tx.upload_blob(&bytes, &PackKey::from_hash(key)).unwrap();
}
#[test]
fn probe_chain_depth_counts_nodes_and_matches_resolve_pack_chain() {
let tx = MemoryTransport::new();
let n1 = h("n1");
let n2 = h("n2");
let n3 = h("n3");
put_node(&tx, n1, None, &[h("pack1")]);
put_node(&tx, n2, Some(n1), &[h("pack2")]);
put_node(&tx, n3, Some(n2), &[h("pack3")]);
let probed = probe_chain(&tx, "main", n3).unwrap();
assert_eq!(probed.depth, 3);
let packs = resolve_pack_chain(&tx, "main", n3).unwrap();
assert_eq!(packs, vec![h("pack1"), h("pack2"), h("pack3")]);
assert_eq!(probed.depth, packs.len());
assert_eq!(probed.packs, packs);
assert_eq!(probed.head, n3);
}
#[test]
fn probe_chain_depth_of_a_single_node_chain_is_one() {
let tx = MemoryTransport::new();
let solo = h("solo");
put_node(&tx, solo, None, &[h("pack-solo")]);
assert_eq!(probe_chain(&tx, "main", solo).unwrap().depth, 1);
}
#[test]
fn probe_chain_errors_on_a_cycle_exactly_like_resolve_pack_chain() {
let tx = MemoryTransport::new();
let a = h("cycle-a");
let b = h("cycle-b");
put_node(&tx, a, Some(b), &[h("pack-a")]);
put_node(&tx, b, Some(a), &[h("pack-b")]);
assert!(matches!(
probe_chain(&tx, "main", a).unwrap_err(),
DispatchError::PackChainInvalid { .. }
));
assert!(matches!(
resolve_pack_chain(&tx, "main", a).unwrap_err(),
DispatchError::PackChainInvalid { .. }
));
}
#[test]
fn probe_chain_errors_on_an_undownloadable_node_like_resolve_pack_chain() {
let tx = MemoryTransport::new();
let ghost = h("never-uploaded");
assert!(matches!(
probe_chain(&tx, "main", ghost).unwrap_err(),
DispatchError::PackChainInvalid { .. }
));
assert!(matches!(
resolve_pack_chain(&tx, "main", ghost).unwrap_err(),
DispatchError::PackChainInvalid { .. }
));
}
fn decode_node_at(tx: &MemoryTransport, key: Hash) -> transfer::PackListNode {
let bytes = tx.download_blob(&PackKey::from_hash(key)).unwrap();
transfer::decode_packlist(&bytes).unwrap()
}
#[test]
fn advance_packmap_multi_key_first_push_writes_one_node_in_order() {
let tx = MemoryTransport::new();
let (k1, k2, k3) = (h("k1"), h("k2"), h("k3"));
let tip = h("tip");
advance_packmap(
&tx,
"main",
&[k1, k2, k3],
ChainAction::Append {
self_contained: true,
},
None,
refs::RefWriteCondition::Missing,
tip,
)
.unwrap();
let pm_head = tx.read_ref(&packmap_ref("main")).unwrap().unwrap();
let node = decode_node_at(&tx, pm_head);
assert_eq!(node.prev, None, "first push has no prior chain");
assert_eq!(
node.packs,
vec![k1, k2, k3],
"all keys land on ONE node, in build/apply order"
);
assert_eq!(tx.read_ref("refs/heads/main").unwrap(), Some(tip));
}
#[test]
fn advance_packmap_single_key_matches_pre_831_shape() {
let tx = MemoryTransport::new();
let k1 = h("only-key");
let tip = h("tip");
advance_packmap(
&tx,
"main",
&[k1],
ChainAction::Append {
self_contained: true,
},
None,
refs::RefWriteCondition::Missing,
tip,
)
.unwrap();
let pm_head = tx.read_ref(&packmap_ref("main")).unwrap().unwrap();
assert_eq!(decode_node_at(&tx, pm_head).packs, vec![k1]);
}
#[test]
fn advance_packmap_is_idempotent_when_every_key_already_chained() {
let tx = MemoryTransport::new();
let (k1, k2, k3) = (h("k1"), h("k2"), h("k3"));
let prior_head = h("prior-node");
put_node(&tx, prior_head, None, &[k1, k2, k3]);
tx.update_ref(
&packmap_ref("main"),
refs::RefWriteCondition::Missing,
&prior_head,
)
.unwrap();
let tip = h("tip");
advance_packmap(
&tx,
"main",
&[k1, k2, k3],
ChainAction::Append {
self_contained: true,
},
None,
refs::RefWriteCondition::Missing,
tip,
)
.unwrap();
assert_eq!(
tx.read_ref(&packmap_ref("main")).unwrap(),
Some(prior_head),
"idempotent retry must not write a new node"
);
assert_eq!(tx.read_ref("refs/heads/main").unwrap(), Some(tip));
}
#[test]
fn advance_packmap_appends_when_only_some_keys_already_chained() {
let tx = MemoryTransport::new();
let (k1, k2, k3) = (h("k1"), h("k2"), h("k3"));
let prior_head = h("prior-node");
put_node(&tx, prior_head, None, &[k1]);
tx.update_ref(
&packmap_ref("main"),
refs::RefWriteCondition::Missing,
&prior_head,
)
.unwrap();
let tip = h("tip");
advance_packmap(
&tx,
"main",
&[k1, k2, k3],
ChainAction::Append {
self_contained: true,
},
None,
refs::RefWriteCondition::Missing,
tip,
)
.unwrap();
let new_head = tx.read_ref(&packmap_ref("main")).unwrap().unwrap();
assert_ne!(new_head, prior_head, "a new node must be appended");
let node = decode_node_at(&tx, new_head);
assert_eq!(node.prev, Some(prior_head));
assert_eq!(node.packs, vec![k1, k2, k3]);
assert_eq!(
resolve_pack_chain(&tx, "main", new_head).unwrap(),
vec![k1, k1, k2, k3],
);
}
}