Skip to main content

mkit_server/pipeline/shard/
d34.rs

1//! Fixed D34 routing: one branch shard, a namespace coordinator, and
2//! enumerable membership and ref-name index shards.
3
4use 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/// The D34 shard map. Fan-outs and head/packmap co-location are permanent
13/// storage contracts; changing them would strand existing rows.
14#[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}