mkit_server/pipeline/shard/
d34.rs1use mkit_core::hash::Hash;
5use mkit_core::hash::hash;
6use mkit_core::refs::PACKMAP_REF_PREFIX;
7
8use super::ShardMap;
9use crate::repo::{NamespaceKey, RepoId};
10use crate::store::{BlobKey, Partition, REF_INDEX_FANOUT};
11
12#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
15pub struct D34Shards;
16
17impl ShardMap for D34Shards {
18 fn ref_shard(&self, repo: &RepoId, ref_name: &str) -> Partition {
19 let shard_ref = ref_name.strip_prefix(PACKMAP_REF_PREFIX).map_or_else(
20 || ref_name.to_owned(),
21 |branch| format!("refs/heads/{branch}"),
22 );
23 Partition::Ref {
24 ns: repo.namespace.clone(),
25 repo: repo.name.clone(),
26 shard_ref,
27 }
28 }
29
30 fn coordinator(&self, ns: &NamespaceKey) -> Partition {
31 Partition::Coordinator(ns.clone())
32 }
33
34 fn ref_index(&self, repo: &RepoId, ref_name: &str) -> Partition {
35 let digest = hash(ref_name.as_bytes());
36 Partition::RefIndex {
37 ns: repo.namespace.clone(),
38 repo: repo.name.clone(),
39 bucket: u16::from_be_bytes([digest[0], digest[1]]) % REF_INDEX_FANOUT,
40 }
41 }
42
43 fn ref_index_partitions(&self, repo: &RepoId) -> Vec<Partition> {
44 (0..REF_INDEX_FANOUT)
45 .map(|bucket| Partition::RefIndex {
46 ns: repo.namespace.clone(),
47 repo: repo.name.clone(),
48 bucket,
49 })
50 .collect()
51 }
52
53 fn membership(&self, repo: &RepoId, pack: &BlobKey) -> Partition {
54 let p = pack.hash();
55 Partition::RepoIndex {
56 ns: repo.namespace.clone(),
57 repo: repo.name.clone(),
58 prefix: (u16::from(p[0]) << 4) | u16::from(p[1] >> 4),
59 }
60 }
61
62 fn object_index(&self, repo: &RepoId, object: &Hash) -> Partition {
63 Partition::RepoIndex {
64 ns: repo.namespace.clone(),
65 repo: repo.name.clone(),
66 prefix: (u16::from(object[0]) << 4) | u16::from(object[1] >> 4),
67 }
68 }
69}
70
71#[cfg(test)]
72mod tests {
73 use std::collections::BTreeSet;
74 use std::path::PathBuf;
75
76 use mkit_core::hash::{to_hex, to_hex_bytes};
77 use proptest::prelude::*;
78 use serde::Serialize;
79
80 use super::*;
81 use crate::repo::RepoName;
82 use crate::store::INDEX_FANOUT;
83
84 fn repo(ns: &str, name: &str) -> RepoId {
85 RepoId {
86 namespace: NamespaceKey::from_stored(ns.into()),
87 name: RepoName::new(name).unwrap(),
88 }
89 }
90
91 fn encoded(partition: &Partition) -> String {
92 let bytes = partition.encode().unwrap();
93 assert_eq!(Partition::decode(&bytes).unwrap(), *partition);
94 to_hex_bytes(&bytes)
95 }
96
97 #[derive(Serialize)]
98 struct RefMapping {
99 namespace: String,
100 repo: String,
101 ref_name: String,
102 ref_partition_hex: String,
103 ref_index_partition_hex: String,
104 }
105
106 #[derive(Serialize)]
107 struct PackMapping {
108 namespace: String,
109 repo: String,
110 pack_id: String,
111 prefix: u16,
112 membership_partition_hex: String,
113 object_index_partition_hex: String,
114 }
115
116 #[derive(Serialize)]
117 struct Mapping {
118 ref_index_fanout: u16,
119 membership_fanout: u16,
120 refs: Vec<RefMapping>,
121 packs: Vec<PackMapping>,
122 coordinator_partition_hex: String,
123 ref_index_partitions_hex: Vec<String>,
124 }
125
126 fn golden_mapping() -> String {
127 let refs = [
128 ("root", "a", "refs/heads/main"),
129 ("root", "a", "refs/mkit/packmap/main"),
130 ("root", "a", "refs/tags/v1"),
131 ("root", "a", "refs/heads/a/b"),
132 ("root", "a", "refs/mkit/packmap/a/b"),
133 ("root", "b", "refs/heads/main"),
134 ("tenant", "a", "refs/heads/main"),
135 ("tenant", "a", "refs/packmaps/main"),
136 ]
137 .into_iter()
138 .map(|(ns, name, ref_name)| {
139 let r = repo(ns, name);
140 RefMapping {
141 namespace: ns.into(),
142 repo: name.into(),
143 ref_name: ref_name.into(),
144 ref_partition_hex: encoded(&D34Shards.ref_shard(&r, ref_name)),
145 ref_index_partition_hex: encoded(&D34Shards.ref_index(&r, ref_name)),
146 }
147 })
148 .collect();
149 let r = repo("root", "a");
150 let packs = [
151 (0x00, 0x00, 0x000),
152 (0xff, 0xff, 0xfff),
153 (0x12, 0x3f, 0x123),
154 ]
155 .into_iter()
156 .map(|(first, second, prefix)| {
157 let mut id = [0; 32];
158 id[..2].copy_from_slice(&[first, second]);
159 let pack = BlobKey::pack(id);
160 PackMapping {
161 namespace: "root".into(),
162 repo: "a".into(),
163 pack_id: to_hex(&id),
164 prefix,
165 membership_partition_hex: encoded(&D34Shards.membership(&r, &pack)),
166 object_index_partition_hex: encoded(&D34Shards.object_index(&r, &id)),
167 }
168 })
169 .collect();
170 let mapping = Mapping {
171 ref_index_fanout: REF_INDEX_FANOUT,
172 membership_fanout: INDEX_FANOUT,
173 refs,
174 packs,
175 coordinator_partition_hex: encoded(&D34Shards.coordinator(&r.namespace)),
176 ref_index_partitions_hex: D34Shards
177 .ref_index_partitions(&r)
178 .iter()
179 .map(encoded)
180 .collect(),
181 };
182 format!("{}\n", serde_json::to_string_pretty(&mapping).unwrap())
183 }
184
185 #[test]
186 fn d34_mapping_golden() {
187 let dir = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../../tests/golden/shards");
188 let mapping = golden_mapping();
189 let manifest = format!(
190 "# WP-1.22 D34 shard mapping. <file> <blake3>\n\
191 # Regenerate: UPDATE_GOLDEN=1 cargo test -p mkit-server d34_mapping_golden\n\
192 d34-mapping.json {}\n",
193 to_hex(&hash(mapping.as_bytes())),
194 );
195 if std::env::var("UPDATE_GOLDEN").as_deref() == Ok("1") {
196 std::fs::create_dir_all(&dir).unwrap();
197 std::fs::write(dir.join("d34-mapping.json"), &mapping).unwrap();
198 std::fs::write(dir.join("MANIFEST.txt"), &manifest).unwrap();
199 return;
200 }
201 assert_eq!(
202 std::fs::read_to_string(dir.join("d34-mapping.json")).unwrap(),
203 mapping
204 );
205 assert_eq!(
206 std::fs::read_to_string(dir.join("MANIFEST.txt")).unwrap(),
207 manifest
208 );
209 }
210
211 #[test]
212 fn d34_ref_index_partitions_are_distinct_and_ordered() {
213 let r = repo("root", "a");
214 let partitions = D34Shards.ref_index_partitions(&r);
215 assert_eq!(partitions.len(), usize::from(REF_INDEX_FANOUT));
216 assert_eq!(
217 partitions.iter().collect::<BTreeSet<_>>().len(),
218 partitions.len()
219 );
220 for (bucket, partition) in (0..REF_INDEX_FANOUT).zip(partitions) {
221 assert_eq!(
222 partition,
223 Partition::RefIndex {
224 ns: r.namespace.clone(),
225 repo: r.name.clone(),
226 bucket,
227 }
228 );
229 }
230 }
231
232 proptest! {
233 #[test]
234 fn d34_branch_and_packmap_share_the_ref_shard(
235 branch in "[a-zA-Z0-9_-]{1,16}(/[a-zA-Z0-9_-]{1,16}){0,4}",
236 ) {
237 let r = repo("root", "a");
238 let head = format!("refs/heads/{branch}");
239 let packmap = format!("{PACKMAP_REF_PREFIX}{branch}");
240 prop_assert_eq!(D34Shards.ref_shard(&r, &head), D34Shards.ref_shard(&r, &packmap));
241 }
242
243 #[test]
244 fn d34_index_buckets_and_membership_prefixes_stay_in_fixed_ranges(
245 name in ".*",
246 id in any::<[u8; 32]>(),
247 ) {
248 let r = repo("root", "a");
249 let Partition::RefIndex { bucket, .. } = D34Shards.ref_index(&r, &name) else {
250 unreachable!("D34 ref-name index is a RefIndex");
251 };
252 let Partition::RepoIndex { prefix, .. } = D34Shards.membership(&r, &BlobKey::pack(id)) else {
253 unreachable!("D34 membership index is a RepoIndex");
254 };
255 prop_assert!(bucket < REF_INDEX_FANOUT);
256 prop_assert!(prefix < INDEX_FANOUT);
257 prop_assert_eq!(prefix, u16::from_be_bytes([id[0], id[1]]) >> 4);
258 prop_assert_eq!(D34Shards.object_index(&r, &id), D34Shards.membership(&r, &BlobKey::pack(id)));
259 }
260 }
261}