1use bytes::{BufMut, Bytes, BytesMut};
91use mkit_core::hash::Hash;
92
93use super::Partition;
94use super::error::StoreError;
95use super::kv::{Key, MAX_KEY_BYTES};
96use crate::quota::QuotaScope;
97use crate::refs::{MAX_REF_NAME_BYTES, is_served_ref_name};
98use crate::repo::{MAX_REPO_NAME_BYTES, NamespaceKey, RepoName};
99
100const _: () = assert!(3 + MAX_REPO_NAME_BYTES + 1 + MAX_REF_NAME_BYTES + 1 + 64 <= MAX_KEY_BYTES);
102
103const _: () = assert!(3 + 1 + 72 + MAX_REPO_NAME_BYTES + MAX_REF_NAME_BYTES + 3 <= MAX_KEY_BYTES);
105
106const _: () = assert!(2 + MAX_REPO_NAME_BYTES + 1 + MAX_REF_NAME_BYTES <= MAX_KEY_BYTES);
108
109pub const LAYOUT_VERSION: u32 = 1;
112
113pub const TAG_LAYOUT_VERSION: &str = "v";
115pub const TAG_SHARDING_MARKER: &str = "sm";
117pub const TAG_ADDRESSING_MARKER: &str = "am";
119pub const TAG_INSPECTION_MARKER: &str = "im";
121pub const TAG_INSPECTION_FLAG: &str = "if";
123pub const TAG_INSPECTION_VERSION: &str = "iv";
125pub const TAG_INSPECTION_HOLD: &str = "ih";
127pub const TAG_INSPECTION_HOLD_INDEX: &str = "ia";
129pub const TAG_INSPECTION_HOLD_MANIFEST: &str = "ir";
131pub const TAG_REF: &str = "r";
133pub const TAG_REF_INDEX: &str = "x";
135pub const TAG_REPLAY: &str = "p";
137pub const TAG_REPLAY_EXPIRY: &str = "px";
139pub const TAG_QUOTA: &str = "q";
141pub const TAG_QUOTA_WINDOW: &str = "qx";
143pub const TAG_QUOTA_SHARD: &str = "qs";
145pub const TAG_QUOTA_VIEW: &str = "qv";
147pub const TAG_QUOTA_CONTRIBUTION: &str = "qc";
149pub const TAG_QUOTA_TOTAL: &str = "qt";
151pub const TAG_GRANT_EPOCH: &str = "e";
153pub const TAG_AUTHORITY_GENERATION: &str = "ag";
155pub const TAG_FENCE_CURSOR: &str = "fc";
157pub const TAG_EPOCH_LEASE: &str = "el";
159pub const TAG_LEASED_SHARD: &str = "ls";
161pub const TAG_LEASE_RECOVERY: &str = "lr";
163pub const TAG_LEASE_RECONCILE: &str = "lrc";
165pub const TAG_BACKUP_STATE: &str = "bk";
167pub const TAG_TIMER: &str = "w";
169pub const TAG_HOLDER: &str = "h";
171pub const TAG_HOLD: &str = "g";
173pub const TAG_PENDING_HOLDER: &str = "gp";
175pub const TAG_CONTENT_TAKEDOWN: &str = "ct";
177pub const TAG_BLOCK: &str = "b";
179pub const TAG_OBJECT_STATE: &str = "c";
181
182pub const TAG_NAMESPACE_RECORD: &str = "nr";
185pub const TAG_REPO_KNOWN: &str = "rk";
187pub const TAG_REPO_REGISTRY: &str = "rr";
190pub const TAG_REPO_LIST: &str = "rl";
192pub const TAG_REPO_VISIBILITY: &str = "rv";
195pub const TAG_NAMESPACE_LIST: &str = "nl";
199
200pub const TAG_TICKET: &str = "t";
202pub const TAG_TICKET_INDEX: &str = "ti";
204pub const TAG_TICKETS_PER_REF: &str = "tc";
206pub const TAG_TICKETS_PER_SIGNER: &str = "tu";
208pub const TAG_MEMBERSHIP: &str = "m";
210pub const TAG_VERIFICATION: &str = "vs";
212pub const TAG_VERIFY_CURSOR: &str = "vc";
214pub const TAG_OBJECT_INDEX: &str = "i";
216pub const TAG_RESERVATION: &str = "o";
218pub const TAG_OUTCOME_PENDING: &str = "oq";
220pub const TAG_RELAY: &str = "or";
222pub const TAG_RELAY_HIGH_WATER: &str = "rh";
224pub const TAG_RELAY_SCAN: &str = "rs";
226pub const TAG_OUTBOX_SEQUENCE: &str = "os";
228pub const TAG_OUTCOME_BACKLOG: &str = "oc";
230pub const TAG_CACHE_PURGE: &str = "cp";
232pub const TAG_CACHE_PURGE_GENERATION: &str = "cg";
234
235pub fn cache_purge(id: &str) -> Result<Key, StoreError> {
237 if id.is_empty()
238 || id.len() > 128
239 || !id
240 .bytes()
241 .all(|b| b.is_ascii_alphanumeric() || b"._:-".contains(&b))
242 {
243 return Err(StoreError::Invalid("invalid purge id".into()));
244 }
245 Ok(key(TAG_CACHE_PURGE, &[id.as_bytes()]))
246}
247#[must_use]
249pub fn cache_purge_generation(scope: &str) -> Key {
250 key(TAG_CACHE_PURGE_GENERATION, &[scope.as_bytes()])
251}
252
253pub const TAG_PUBLICATION: &str = "pp";
255pub const TAG_ADVANCE: &str = "av";
257pub const TAG_PUBLISHED_REF: &str = "pr";
259pub const TAG_PUBLISHED_INDEX: &str = "py";
261pub const TAG_PUBLISHED_MEMBER: &str = "pm";
263
264pub const RESERVED_TAGS: &[&str] = &["tb", "l", TAG_NAMESPACE_LIST];
266
267pub const VC_JOB: u8 = 0;
269pub const VC_FRAME: u8 = 1;
271pub const VC_CHILD: u8 = 2;
273pub const VC_BASE: u8 = 3;
275pub const VC_CANDIDATE: u8 = 4;
277pub const VC_HISTORY: u8 = 5;
280pub const VC_DEPENDENCY: u8 = 6;
282
283#[derive(Debug, Clone, PartialEq, Eq)]
285#[non_exhaustive]
286pub enum ParsedKey {
287 CachePurge(String),
289 CachePurgeGeneration(String),
291 BackupState,
293 ShardingMarker,
295 AddressingMarker,
297 InspectionMarker,
299 InspectionFlag { repo: RepoName, id: Hash },
301 InspectionVersion(RepoName),
303 InspectionHold {
305 repo: RepoName,
306 content: Hash,
307 advance: Hash,
308 },
309 InspectionHoldIndex { repo: RepoName, advance: Hash },
311 InspectionHoldManifest { repo: RepoName, advance: Hash },
313 LayoutVersion,
315 NamespaceRecord,
317 RepoRecord(RepoName),
319 RepoListing {
321 explicit_public: bool,
322 repo: RepoName,
323 },
324 RepoVisibility(RepoName),
326 RepoKnown(RepoName),
328 RelayHighWater(Partition),
330 Ref {
332 repo: RepoName,
334 name: String,
336 },
337 RefIndexEntry {
339 repo: RepoName,
341 name: String,
343 },
344 Publication { repo: RepoName, name: String },
346 Advance {
348 repo: RepoName,
349 name: String,
350 sequence: u64,
351 },
352 PublishedRef { repo: RepoName, name: String },
354 PublishedIndex { repo: RepoName, name: String },
356 PublishedMember { repo: RepoName, pack_id: Hash },
358 Replay(Hash),
360 ReplayExpiry {
362 expires_at_ms: u64,
364 scope: Hash,
366 },
367 Quota(String),
369 QuotaWindow {
371 window_start_ms: u64,
373 scope: String,
375 },
376 QuotaShard(u64),
378 QuotaView(u64),
380 QuotaContribution { window: u64, source: Partition },
382 QuotaTotal(u64),
384 Ticket(Hash),
386 TicketIndex {
388 repo: RepoName,
390 name: String,
392 pack_id: Hash,
394 signer: Hash,
396 },
397 TicketsPerRef {
399 repo: RepoName,
401 name: String,
403 },
404 TicketsPerSigner {
406 repo: RepoName,
408 name: String,
410 signer: Hash,
412 },
413 Membership {
415 repo: RepoName,
417 pack_id: Hash,
419 },
420 Verification { repo: RepoName, pack_id: Hash },
422 VerifyCursor {
424 repo: RepoName,
426 pack_id: Hash,
428 sub: u8,
430 id: Option<Hash>,
432 },
433 ObjectIndex {
435 repo: RepoName,
437 object: Hash,
439 pack_id: Hash,
441 },
442 Reservation(String),
444 OutcomePending {
446 seq: u64,
448 reservation_id: String,
450 },
451 Relay(u64),
453 RelayScan,
455 OutboxSequence,
457 OutcomeBacklog,
459 GrantEpoch,
461 AuthorityGeneration,
463 EpochLease,
465 LeasedShard {
467 repo: RepoName,
469 shard_ref: String,
471 },
472 LeaseRecovery,
474 FenceCursor(u8),
476 LeaseReconcile,
478 Timer {
480 due_at_ms: u64,
482 kind: u8,
484 reference: Bytes,
486 },
487 Holder {
489 object: Hash,
491 ns: NamespaceKey,
493 repo: RepoName,
495 },
496 Hold {
498 object: Hash,
500 hold_id: Hash,
502 },
503 PendingHolder {
505 object: Hash,
507 hold_id: Hash,
509 },
510 ContentTakedown {
512 object: Hash,
514 intent: Hash,
516 },
517 Block(Hash),
519 ObjectState(Hash),
521}
522
523fn key(tag: &str, parts: &[&[u8]]) -> Key {
524 let mut buf =
525 BytesMut::with_capacity(tag.len() + 1 + parts.iter().map(|p| p.len()).sum::<usize>());
526 buf.put_slice(tag.as_bytes());
527 buf.put_u8(0);
528 for part in parts {
529 buf.put_slice(part);
530 }
531 Key::new(buf.freeze())
532}
533
534fn successor(prefix: &Key) -> Key {
536 let mut bytes = prefix.as_bytes().to_vec();
537 while bytes.last() == Some(&0xff) {
538 bytes.pop();
539 }
540 if let Some(last) = bytes.last_mut() {
541 *last += 1;
542 }
543 Key::new(bytes)
544}
545
546#[must_use]
548pub fn class_range(tag: &str) -> (Key, Key) {
549 let start = key(tag, &[]);
550 let end = successor(&start);
551 (start, end)
552}
553
554#[must_use]
556pub fn is_ref_key(key: &Key) -> bool {
557 key.as_bytes().starts_with(b"r\0")
558}
559
560#[must_use]
562pub fn layout_version() -> Key {
563 key(TAG_LAYOUT_VERSION, &[])
564}
565
566#[must_use]
568pub fn sharding_marker() -> Key {
569 key(TAG_SHARDING_MARKER, &[])
570}
571
572#[must_use]
574pub fn addressing_marker() -> Key {
575 key(TAG_ADDRESSING_MARKER, &[])
576}
577
578#[must_use]
580pub fn inspection_marker() -> Key {
581 key(TAG_INSPECTION_MARKER, &[])
582}
583
584#[must_use]
586pub fn inspection_flag(repo: &RepoName, id: &Hash) -> Key {
587 key(TAG_INSPECTION_FLAG, &[repo.as_str().as_bytes(), &[0], id])
588}
589
590#[must_use]
592pub fn inspection_version(repo: &RepoName) -> Key {
593 key(TAG_INSPECTION_VERSION, &[repo.as_str().as_bytes()])
594}
595
596#[must_use]
598pub fn inspection_hold(repo: &RepoName, content: &Hash, advance: &Hash) -> Key {
599 key(
600 TAG_INSPECTION_HOLD,
601 &[repo.as_str().as_bytes(), &[0], content, advance],
602 )
603}
604
605#[must_use]
607pub fn inspection_hold_range(repo: &RepoName, content: &Hash) -> (Key, Key) {
608 let start = key(
609 TAG_INSPECTION_HOLD,
610 &[repo.as_str().as_bytes(), &[0], content],
611 );
612 let end = successor(&start);
613 (start, end)
614}
615
616#[must_use]
618pub fn inspection_hold_index(repo: &RepoName, advance: &Hash) -> Key {
619 key(
620 TAG_INSPECTION_HOLD_INDEX,
621 &[repo.as_str().as_bytes(), &[0], advance],
622 )
623}
624
625#[must_use]
627pub fn inspection_hold_manifest(repo: &RepoName, advance: &Hash) -> Key {
628 key(
629 TAG_INSPECTION_HOLD_MANIFEST,
630 &[repo.as_str().as_bytes(), &[0], advance],
631 )
632}
633
634#[must_use]
636pub fn namespace_record() -> Key {
637 key(TAG_NAMESPACE_RECORD, &[])
638}
639
640#[must_use]
642pub fn repo_record(repo: &RepoName) -> Key {
643 key(TAG_REPO_REGISTRY, &[repo.as_str().as_bytes()])
644}
645
646#[must_use]
648pub fn repo_listing_prefix(explicit_public: bool) -> Key {
649 key(
650 TAG_REPO_LIST,
651 &[if explicit_public { b"p" } else { b"d" }, &[0]],
652 )
653}
654
655#[must_use]
657pub fn repo_listing(repo: &RepoName, explicit_public: bool) -> Key {
658 let mut bytes = repo_listing_prefix(explicit_public).as_bytes().to_vec();
659 bytes.extend_from_slice(repo.as_str().as_bytes());
660 Key::new(bytes)
661}
662
663#[must_use]
665pub fn repo_visibility(repo: &RepoName) -> Key {
666 key(TAG_REPO_VISIBILITY, &[repo.as_str().as_bytes()])
667}
668
669#[must_use]
671pub fn repo_known(repo: &RepoName) -> Key {
672 key(TAG_REPO_KNOWN, &[repo.as_str().as_bytes()])
673}
674
675#[must_use]
677pub fn ref_key(repo: &RepoName, name: &str) -> Key {
678 key(TAG_REF, &[repo.as_str().as_bytes(), b"\0", name.as_bytes()])
679}
680
681#[must_use]
683pub fn ref_prefix_range(repo: &RepoName, prefix: &str) -> (Key, Key) {
684 let start = ref_key(repo, prefix);
685 let end = successor(&start);
686 (start, end)
687}
688
689#[must_use]
691pub fn ref_index_key(repo: &RepoName, name: &str) -> Key {
692 key(
693 TAG_REF_INDEX,
694 &[repo.as_str().as_bytes(), b"\0", name.as_bytes()],
695 )
696}
697
698#[must_use]
700pub fn ref_index_prefix_range(repo: &RepoName, prefix: &str) -> (Key, Key) {
701 let start = ref_index_key(repo, prefix);
702 let end = successor(&start);
703 (start, end)
704}
705
706#[must_use]
708pub fn publication(repo: &RepoName, name: &str) -> Key {
709 key(
710 TAG_PUBLICATION,
711 &[repo.as_str().as_bytes(), b"\0", name.as_bytes()],
712 )
713}
714#[must_use]
716pub fn advance(repo: &RepoName, name: &str, sequence: u64) -> Key {
717 key(
718 TAG_ADVANCE,
719 &[
720 repo.as_str().as_bytes(),
721 b"\0",
722 name.as_bytes(),
723 b"\0",
724 &sequence.to_be_bytes(),
725 ],
726 )
727}
728#[must_use]
730pub fn published_ref(repo: &RepoName, name: &str) -> Key {
731 key(
732 TAG_PUBLISHED_REF,
733 &[repo.as_str().as_bytes(), b"\0", name.as_bytes()],
734 )
735}
736#[must_use]
738pub fn published_index(repo: &RepoName, name: &str) -> Key {
739 key(
740 TAG_PUBLISHED_INDEX,
741 &[repo.as_str().as_bytes(), b"\0", name.as_bytes()],
742 )
743}
744#[must_use]
746pub fn published_member(repo: &RepoName, pack: &Hash) -> Key {
747 key(
748 TAG_PUBLISHED_MEMBER,
749 &[repo.as_str().as_bytes(), b"\0", pack],
750 )
751}
752#[must_use]
754pub fn published_range(repo: &RepoName, prefix: &str, index: bool) -> (Key, Key) {
755 let start = if index {
756 published_index(repo, prefix)
757 } else {
758 published_ref(repo, prefix)
759 };
760 let end = successor(&start);
761 (start, end)
762}
763
764#[must_use]
766pub fn replay(scope: &Hash) -> Key {
767 key(TAG_REPLAY, &[scope])
768}
769
770#[must_use]
772pub fn replay_expiry(expires_at_ms: u64, scope: &Hash) -> Key {
773 key(TAG_REPLAY_EXPIRY, &[&expires_at_ms.to_be_bytes(), scope])
774}
775
776#[must_use]
778pub fn replay_expiry_before(before_ms: u64) -> (Key, Key) {
779 let (start, _) = class_range(TAG_REPLAY_EXPIRY);
780 (start, key(TAG_REPLAY_EXPIRY, &[&before_ms.to_be_bytes()]))
781}
782
783#[must_use]
785pub fn quota(scope: &QuotaScope) -> Key {
786 key(TAG_QUOTA, &[scope.as_str().as_bytes()])
787}
788
789#[must_use]
791pub fn quota_window(window_start_ms: u64, scope: &QuotaScope) -> Key {
792 key(
793 TAG_QUOTA_WINDOW,
794 &[&window_start_ms.to_be_bytes(), scope.as_str().as_bytes()],
795 )
796}
797
798#[must_use]
800pub fn quota_window_before(before_ms: u64) -> (Key, Key) {
801 let (start, _) = class_range(TAG_QUOTA_WINDOW);
802 (start, key(TAG_QUOTA_WINDOW, &[&before_ms.to_be_bytes()]))
803}
804
805#[must_use]
807pub fn quota_shard(window: u64) -> Key {
808 key(TAG_QUOTA_SHARD, &[&window.to_be_bytes()])
809}
810
811#[must_use]
813pub fn quota_view(window: u64) -> Key {
814 key(TAG_QUOTA_VIEW, &[&window.to_be_bytes()])
815}
816
817pub fn quota_contribution(window: u64, source: &Partition) -> Result<Key, StoreError> {
819 if !matches!(source, Partition::Ref { .. }) {
820 return Err(StoreError::Invalid(
821 "quota source is not a ref shard".into(),
822 ));
823 }
824 checked_key(
825 TAG_QUOTA_CONTRIBUTION,
826 &[&window.to_be_bytes(), &source.encode()?],
827 )
828}
829
830#[must_use]
832pub fn quota_total(window: u64) -> Key {
833 key(TAG_QUOTA_TOTAL, &[&window.to_be_bytes()])
834}
835
836#[must_use]
838pub fn quota_namespace_before(tag: &str, window: u64) -> (Key, Key) {
839 (key(tag, &[]), key(tag, &[&window.to_be_bytes()]))
840}
841
842#[must_use]
844pub fn quota_namespace_window(tag: &str, window: u64) -> (Key, Key) {
845 (
846 key(tag, &[&window.to_be_bytes()]),
847 key(tag, &[&window.saturating_add(1).to_be_bytes()]),
848 )
849}
850
851#[must_use]
853pub fn validate_reservation_id(rid: &str) -> bool {
854 (1..=128).contains(&rid.len())
855 && rid
856 .bytes()
857 .all(|b| b.is_ascii_alphanumeric() || b"._:-".contains(&b))
858}
859
860fn checked_key(tag: &str, parts: &[&[u8]]) -> Result<Key, StoreError> {
861 let result = key(tag, parts);
862 if result.as_bytes().len() > MAX_KEY_BYTES {
863 return Err(StoreError::Invalid("key exceeds MAX_KEY_BYTES".into()));
864 }
865 Ok(result)
866}
867
868fn check_ticket_ref(name: &str) -> Result<(), StoreError> {
869 if !is_served_ref_name(name) {
870 return Err(StoreError::Invalid("invalid ticket ref name".into()));
871 }
872 Ok(())
873}
874
875#[must_use]
877pub fn ticket(id: &Hash) -> Key {
878 key(TAG_TICKET, &[id])
879}
880
881pub fn ticket_index(
883 repo: &RepoName,
884 name: &str,
885 pack: &Hash,
886 signer: &Hash,
887) -> Result<Key, StoreError> {
888 check_ticket_ref(name)?;
889 checked_key(
890 TAG_TICKET_INDEX,
891 &[
892 repo.as_str().as_bytes(),
893 b"\0",
894 name.as_bytes(),
895 b"\0",
896 pack,
897 signer,
898 ],
899 )
900}
901
902pub fn tickets_per_ref(repo: &RepoName, name: &str) -> Result<Key, StoreError> {
904 check_ticket_ref(name)?;
905 checked_key(
906 TAG_TICKETS_PER_REF,
907 &[repo.as_str().as_bytes(), b"\0", name.as_bytes()],
908 )
909}
910
911pub fn tickets_per_signer(repo: &RepoName, name: &str, signer: &Hash) -> Result<Key, StoreError> {
913 check_ticket_ref(name)?;
914 checked_key(
915 TAG_TICKETS_PER_SIGNER,
916 &[
917 repo.as_str().as_bytes(),
918 b"\0",
919 name.as_bytes(),
920 b"\0",
921 signer,
922 ],
923 )
924}
925
926#[must_use]
928pub fn membership(repo: &RepoName, pack: &Hash) -> Key {
929 key(TAG_MEMBERSHIP, &[repo.as_str().as_bytes(), b"\0", pack])
930}
931
932#[must_use]
934pub fn verification(repo: &RepoName, pack: &Hash) -> Key {
935 key(TAG_VERIFICATION, &[repo.as_str().as_bytes(), b"\0", pack])
936}
937
938#[must_use]
940pub fn verify_job(repo: &RepoName, pack: &Hash) -> Key {
941 verify_row(repo, pack, VC_JOB, None)
942}
943
944#[must_use]
946pub fn verify_row(repo: &RepoName, pack: &Hash, sub: u8, id: Option<&Hash>) -> Key {
947 key(
948 TAG_VERIFY_CURSOR,
949 &[
950 repo.as_str().as_bytes(),
951 b"\0",
952 pack,
953 &[sub],
954 id.map_or(&[][..], |id| &id[..]),
955 ],
956 )
957}
958
959#[must_use]
961pub fn verify_range(repo: &RepoName, pack: &Hash, sub: Option<u8>) -> (Key, Key) {
962 let mut parts: Vec<&[u8]> = vec![repo.as_str().as_bytes(), b"\0", pack];
963 let sub = sub.map(|sub| [sub]);
964 if let Some(sub) = &sub {
965 parts.push(sub);
966 }
967 let start = key(TAG_VERIFY_CURSOR, &parts);
968 let end = successor(&start);
969 (start, end)
970}
971
972#[must_use]
974pub fn object_index(repo: &RepoName, object: &Hash, pack: &Hash) -> Key {
975 key(
976 TAG_OBJECT_INDEX,
977 &[repo.as_str().as_bytes(), b"\0", object, pack],
978 )
979}
980
981#[must_use]
983pub fn object_index_range(repo: &RepoName, object: &Hash) -> (Key, Key) {
984 let start = key(TAG_OBJECT_INDEX, &[repo.as_str().as_bytes(), b"\0", object]);
985 let end = successor(&start);
986 (start, end)
987}
988
989pub fn reservation(rid: &str) -> Result<Key, StoreError> {
991 if !validate_reservation_id(rid) {
992 return Err(StoreError::Invalid("invalid reservation id".into()));
993 }
994 checked_key(TAG_RESERVATION, &[rid.as_bytes()])
995}
996
997pub fn outcome_pending(seq: u64, rid: &str) -> Result<Key, StoreError> {
999 if !validate_reservation_id(rid) || seq == 0 {
1000 return Err(StoreError::Invalid("invalid outcome pending key".into()));
1001 }
1002 checked_key(TAG_OUTCOME_PENDING, &[&seq.to_be_bytes(), rid.as_bytes()])
1003}
1004
1005#[must_use]
1007pub fn relay(seq: u64) -> Key {
1008 key(TAG_RELAY, &[&seq.to_be_bytes()])
1009}
1010
1011pub fn relay_high_water(source: &Partition) -> Result<Key, StoreError> {
1013 checked_key(TAG_RELAY_HIGH_WATER, &[&source.encode()?])
1014}
1015
1016#[must_use]
1018pub fn relay_scan() -> Key {
1019 key(TAG_RELAY_SCAN, &[])
1020}
1021
1022#[must_use]
1024pub fn outbox_sequence() -> Key {
1025 key(TAG_OUTBOX_SEQUENCE, &[])
1026}
1027
1028#[must_use]
1030pub fn outcome_backlog() -> Key {
1031 key(TAG_OUTCOME_BACKLOG, &[])
1032}
1033
1034#[must_use]
1036pub fn authority_generation() -> Key {
1037 key(TAG_AUTHORITY_GENERATION, &[])
1038}
1039
1040#[must_use]
1042pub fn revoke_cursor(authority: bool) -> Key {
1043 key(TAG_FENCE_CURSOR, &[&[u8::from(authority)]])
1044}
1045
1046#[must_use]
1048pub fn grant_epoch() -> Key {
1049 key(TAG_GRANT_EPOCH, &[])
1050}
1051
1052#[must_use]
1054pub fn epoch_lease() -> Key {
1055 key(TAG_EPOCH_LEASE, &[])
1056}
1057
1058#[must_use]
1060pub fn leased_shard(repo: &RepoName, shard_ref: &str) -> Key {
1061 key(
1062 TAG_LEASED_SHARD,
1063 &[repo.as_str().as_bytes(), b"\0", shard_ref.as_bytes()],
1064 )
1065}
1066
1067#[must_use]
1069pub fn lease_recovery() -> Key {
1070 key(TAG_LEASE_RECOVERY, &[])
1071}
1072
1073#[must_use]
1075pub fn lease_reconcile() -> Key {
1076 key(TAG_LEASE_RECONCILE, &[])
1077}
1078
1079#[must_use]
1081pub fn backup_state() -> Key {
1082 key(TAG_BACKUP_STATE, &[])
1083}
1084
1085pub const MAX_TIMER_RETRY_ATTEMPT: u8 = 8;
1087
1088#[must_use]
1091pub fn timer(due_at_ms: u64, kind: u8, reference: &[u8]) -> Key {
1092 timer_retry(due_at_ms, kind, reference, due_at_ms, 0)
1093}
1094
1095#[must_use]
1098pub fn timer_retry(
1099 due_at_ms: u64,
1100 kind: u8,
1101 reference: &[u8],
1102 original_due_at_ms: u64,
1103 attempt: u8,
1104) -> Key {
1105 key(
1106 TAG_TIMER,
1107 &[
1108 &due_at_ms.to_be_bytes(),
1109 &[kind, attempt],
1110 &original_due_at_ms.to_be_bytes(),
1111 reference,
1112 ],
1113 )
1114}
1115
1116#[must_use]
1119pub fn timer_retry_state(key: &Key) -> Option<(u64, u8)> {
1120 let body = key.as_bytes().strip_prefix(b"w\0")?;
1121 let (_, _, original_due, attempt, _) = timer_parts(body)?;
1122 Some((original_due, attempt))
1123}
1124
1125fn timer_parts(body: &[u8]) -> Option<(u64, u8, u64, u8, &[u8])> {
1126 let (due_at_ms, rest) = be64(body)?;
1127 let (&kind, rest) = rest.split_first()?;
1128 let (&attempt, rest) = rest.split_first()?;
1129 let (original_due, reference) = be64(rest)?;
1130 if attempt > MAX_TIMER_RETRY_ATTEMPT
1131 || original_due > due_at_ms
1132 || attempt == 0 && original_due != due_at_ms
1133 {
1134 return None;
1135 }
1136 Some((due_at_ms, kind, original_due, attempt, reference))
1137}
1138
1139pub fn holder(object: &Hash, ns: &NamespaceKey, repo: &RepoName) -> Result<Key, StoreError> {
1145 if ns.as_str().as_bytes().contains(&0) {
1146 return Err(StoreError::Invalid("namespace contains 0x00".into()));
1147 }
1148 Ok(key(
1149 TAG_HOLDER,
1150 &[
1151 object,
1152 ns.as_str().as_bytes(),
1153 b"\0",
1154 repo.as_str().as_bytes(),
1155 ],
1156 ))
1157}
1158
1159#[must_use]
1161pub fn holders_of(object: &Hash) -> (Key, Key) {
1162 let start = key(TAG_HOLDER, &[object]);
1163 let end = successor(&start);
1164 (start, end)
1165}
1166
1167#[must_use]
1169pub fn hold(object: &Hash, hold_id: &Hash) -> Key {
1170 key(TAG_HOLD, &[object, hold_id])
1171}
1172
1173#[must_use]
1175pub fn holds_of(object: &Hash) -> (Key, Key) {
1176 let start = key(TAG_HOLD, &[object]);
1177 let end = successor(&start);
1178 (start, end)
1179}
1180
1181#[must_use]
1183pub fn pending_holder(object: &Hash, hold_id: &Hash) -> Key {
1184 key(TAG_PENDING_HOLDER, &[object, hold_id])
1185}
1186
1187#[must_use]
1189pub fn pending_holders_of(object: &Hash) -> (Key, Key) {
1190 let start = key(TAG_PENDING_HOLDER, &[object]);
1191 let end = successor(&start);
1192 (start, end)
1193}
1194
1195#[must_use]
1197pub fn content_takedown(object: &Hash, intent: &Hash) -> Key {
1198 key(TAG_CONTENT_TAKEDOWN, &[object, intent])
1199}
1200
1201#[must_use]
1203pub fn content_takedowns_of(object: &Hash) -> (Key, Key) {
1204 let start = key(TAG_CONTENT_TAKEDOWN, &[object]);
1205 let end = successor(&start);
1206 (start, end)
1207}
1208
1209#[must_use]
1211pub fn block(object: &Hash) -> Key {
1212 key(TAG_BLOCK, &[object])
1213}
1214
1215#[must_use]
1217pub fn object_state(object: &Hash) -> Key {
1218 key(TAG_OBJECT_STATE, &[object])
1219}
1220
1221fn be64(bytes: &[u8]) -> Option<(u64, &[u8])> {
1222 let (head, rest) = bytes.split_first_chunk::<8>()?;
1223 Some((u64::from_be_bytes(*head), rest))
1224}
1225
1226fn hash(bytes: &[u8]) -> Option<Hash> {
1227 Hash::try_from(bytes).ok()
1228}
1229
1230fn parse_ticket_binding(tag: &[u8], body: &[u8]) -> Option<ParsedKey> {
1231 let text = |b: &[u8]| String::from_utf8(b.to_vec()).ok();
1232 Some({
1233 let sep = body.iter().position(|&b| b == 0)?;
1234 let repo = RepoName::new(text(&body[..sep])?).ok()?;
1235 let rest = &body[sep + 1..];
1236 if tag == b"tc" {
1237 let name = text(rest)?;
1238 check_ticket_ref(&name).ok()?;
1239 ParsedKey::TicketsPerRef { repo, name }
1240 } else {
1241 let sep = rest.iter().position(|&b| b == 0)?;
1242 let name = text(&rest[..sep])?;
1243 check_ticket_ref(&name).ok()?;
1244 let tail = &rest[sep + 1..];
1245 if tag == b"ti" {
1246 let (pack_id, signer) = tail.split_first_chunk::<32>()?;
1247 ParsedKey::TicketIndex {
1248 repo,
1249 name,
1250 pack_id: *pack_id,
1251 signer: hash(signer)?,
1252 }
1253 } else {
1254 ParsedKey::TicketsPerSigner {
1255 repo,
1256 name,
1257 signer: hash(tail)?,
1258 }
1259 }
1260 }
1261 })
1262}
1263
1264fn parse_reservation_id(bytes: &[u8]) -> Option<String> {
1265 let rid = std::str::from_utf8(bytes).ok()?;
1266 validate_reservation_id(rid).then(|| rid.to_owned())
1267}
1268
1269fn parse_outcome_pending(body: &[u8]) -> Option<ParsedKey> {
1270 let (seq, rest) = be64(body)?;
1271 let reservation_id = parse_reservation_id(rest)?;
1272 if seq == 0 {
1273 return None;
1274 }
1275 Some(ParsedKey::OutcomePending {
1276 seq,
1277 reservation_id,
1278 })
1279}
1280
1281fn parse_leased_shard(body: &[u8]) -> Option<ParsedKey> {
1282 let sep = body.iter().position(|&b| b == 0)?;
1283 Some(ParsedKey::LeasedShard {
1284 repo: RepoName::new(core::str::from_utf8(&body[..sep]).ok()?).ok()?,
1285 shard_ref: core::str::from_utf8(&body[sep + 1..]).ok()?.to_owned(),
1286 })
1287}
1288
1289fn parse_holder(body: &[u8]) -> Option<ParsedKey> {
1290 let (object, rest) = body.split_first_chunk::<32>()?;
1291 let sep = rest.iter().position(|&b| b == 0)?;
1292 Some(ParsedKey::Holder {
1293 object: *object,
1294 ns: NamespaceKey::from_stored(String::from_utf8(rest[..sep].to_vec()).ok()?),
1295 repo: RepoName::new(String::from_utf8(rest[sep + 1..].to_vec()).ok()?).ok()?,
1296 })
1297}
1298
1299fn parse_namespace_quota(tag: &[u8], body: &[u8]) -> Option<ParsedKey> {
1300 let (window, rest) = be64(body)?;
1301 match tag {
1302 b"qs" if rest.is_empty() => Some(ParsedKey::QuotaShard(window)),
1303 b"qv" if rest.is_empty() => Some(ParsedKey::QuotaView(window)),
1304 b"qt" if rest.is_empty() => Some(ParsedKey::QuotaTotal(window)),
1305 b"qc" => {
1306 let source = Partition::decode(rest).ok()?;
1307 matches!(source, Partition::Ref { .. })
1308 .then_some(ParsedKey::QuotaContribution { window, source })
1309 }
1310 _ => None,
1311 }
1312}
1313
1314fn parse_named_ref(body: &[u8]) -> Option<(RepoName, String)> {
1315 let sep = body.iter().position(|&b| b == 0)?;
1316 let repo = RepoName::new(String::from_utf8(body[..sep].to_vec()).ok()?).ok()?;
1317 let name = String::from_utf8(body[sep + 1..].to_vec()).ok()?;
1318 Some((repo, name))
1319}
1320
1321#[must_use]
1324#[allow(clippy::too_many_lines)] pub fn parse(key: &Key) -> Option<ParsedKey> {
1326 let bytes = key.as_bytes();
1327 if bytes.len() > MAX_KEY_BYTES {
1328 return None;
1329 }
1330 let split = bytes.iter().position(|&b| b == 0)?;
1331 let (tag, body) = (&bytes[..split], &bytes[split + 1..]);
1332 let text = |b: &[u8]| String::from_utf8(b.to_vec()).ok();
1333 Some(match tag {
1334 b"cp" => {
1335 let id = text(body)?;
1336 cache_purge(&id).ok()?;
1337 ParsedKey::CachePurge(id)
1338 }
1339 b"cg" if !body.is_empty() => ParsedKey::CachePurgeGeneration(text(body)?),
1340 b"sm" if body.is_empty() => ParsedKey::ShardingMarker,
1341 b"am" if body.is_empty() => ParsedKey::AddressingMarker,
1342 b"im" if body.is_empty() => ParsedKey::InspectionMarker,
1343 b"iv" => ParsedKey::InspectionVersion(RepoName::new(text(body)?).ok()?),
1344 b"if" | b"ih" | b"ia" | b"ir" => {
1345 let sep = body.iter().position(|&b| b == 0)?;
1346 let repo = RepoName::new(text(&body[..sep])?).ok()?;
1347 let suffix = &body[sep + 1..];
1348 match tag {
1349 b"if" => ParsedKey::InspectionFlag {
1350 repo,
1351 id: hash(suffix)?,
1352 },
1353 b"ia" => ParsedKey::InspectionHoldIndex {
1354 repo,
1355 advance: hash(suffix)?,
1356 },
1357 b"ir" => ParsedKey::InspectionHoldManifest {
1358 repo,
1359 advance: hash(suffix)?,
1360 },
1361 _ => {
1362 let (content, advance) = suffix.split_first_chunk::<32>()?;
1363 ParsedKey::InspectionHold {
1364 repo,
1365 content: *content,
1366 advance: hash(advance)?,
1367 }
1368 }
1369 }
1370 }
1371 b"v" if body.is_empty() => ParsedKey::LayoutVersion,
1372 b"e" if body.is_empty() => ParsedKey::GrantEpoch,
1373 b"ag" if body.is_empty() => ParsedKey::AuthorityGeneration,
1374 b"el" if body.is_empty() => ParsedKey::EpochLease,
1375 b"lr" if body.is_empty() => ParsedKey::LeaseRecovery,
1376 b"fc" if body.len() == 1 && body[0] <= 1 => ParsedKey::FenceCursor(body[0]),
1377 b"lrc" if body.is_empty() => ParsedKey::LeaseReconcile,
1378 b"bk" if body.is_empty() => ParsedKey::BackupState,
1379 b"ls" => parse_leased_shard(body)?,
1380 b"nr" if body.is_empty() => ParsedKey::NamespaceRecord,
1381 b"rr" => ParsedKey::RepoRecord(RepoName::new(text(body)?).ok()?),
1382 b"rl" if body.starts_with(b"p\0") || body.starts_with(b"d\0") => ParsedKey::RepoListing {
1383 explicit_public: body[0] == b'p',
1384 repo: RepoName::new(text(&body[2..])?).ok()?,
1385 },
1386 b"rv" => ParsedKey::RepoVisibility(RepoName::new(text(body)?).ok()?),
1387 b"rh" => ParsedKey::RelayHighWater(Partition::decode(body).ok()?),
1388 b"rs" if body.is_empty() => ParsedKey::RelayScan,
1389 b"rk" => ParsedKey::RepoKnown(RepoName::new(text(body)?).ok()?),
1390 b"r" => {
1391 let (repo, name) = parse_named_ref(body)?;
1392 ParsedKey::Ref { repo, name }
1393 }
1394 b"x" => {
1395 let (repo, name) = parse_named_ref(body)?;
1396 ParsedKey::RefIndexEntry { repo, name }
1397 }
1398 b"pp" | b"pr" | b"py" => {
1399 let (repo, name) = parse_named_ref(body)?;
1400 check_ticket_ref(&name).ok()?;
1401 match tag {
1402 b"pp" => ParsedKey::Publication { repo, name },
1403 b"pr" => ParsedKey::PublishedRef { repo, name },
1404 _ => ParsedKey::PublishedIndex { repo, name },
1405 }
1406 }
1407 b"av" => {
1408 let cut = body.len().checked_sub(9)?;
1409 if body[cut] != 0 {
1410 return None;
1411 }
1412 let (repo, name) = parse_named_ref(&body[..cut])?;
1413 check_ticket_ref(&name).ok()?;
1414 let (sequence, rest) = be64(&body[cut + 1..])?;
1415 if sequence == 0 || !rest.is_empty() {
1416 return None;
1417 }
1418 ParsedKey::Advance {
1419 repo,
1420 name,
1421 sequence,
1422 }
1423 }
1424 b"pm" => {
1425 let sep = body.iter().position(|b| *b == 0)?;
1426 ParsedKey::PublishedMember {
1427 repo: RepoName::new(text(&body[..sep])?).ok()?,
1428 pack_id: hash(&body[sep + 1..])?,
1429 }
1430 }
1431 b"t" => ParsedKey::Ticket(hash(body)?),
1432 b"ti" | b"tc" | b"tu" => parse_ticket_binding(tag, body)?,
1433 b"m" => {
1434 let sep = body.iter().position(|&b| b == 0)?;
1435 ParsedKey::Membership {
1436 repo: RepoName::new(text(&body[..sep])?).ok()?,
1437 pack_id: hash(&body[sep + 1..])?,
1438 }
1439 }
1440 b"vs" => {
1441 let sep = body.iter().position(|&b| b == 0)?;
1442 ParsedKey::Verification {
1443 repo: RepoName::new(text(&body[..sep])?).ok()?,
1444 pack_id: hash(&body[sep + 1..])?,
1445 }
1446 }
1447 b"vc" => {
1448 let sep = body.iter().position(|&b| b == 0)?;
1449 let (pack, rest) = body[sep + 1..].split_first_chunk::<32>()?;
1450 let (sub, id) = rest.split_first()?;
1451 let id = match (*sub, id.len()) {
1452 (VC_JOB, 0) => None,
1453 (VC_FRAME..=VC_DEPENDENCY, 32) => Some(hash(id)?),
1454 _ => return None,
1455 };
1456 ParsedKey::VerifyCursor {
1457 repo: RepoName::new(text(&body[..sep])?).ok()?,
1458 pack_id: *pack,
1459 sub: *sub,
1460 id,
1461 }
1462 }
1463 b"i" => {
1464 let sep = body.iter().position(|&b| b == 0)?;
1465 let (object, pack_id) = body[sep + 1..].split_first_chunk::<32>()?;
1466 ParsedKey::ObjectIndex {
1467 repo: RepoName::new(text(&body[..sep])?).ok()?,
1468 object: *object,
1469 pack_id: hash(pack_id)?,
1470 }
1471 }
1472 b"o" => ParsedKey::Reservation(parse_reservation_id(body)?),
1473 b"oq" => parse_outcome_pending(body)?,
1474 b"or" => {
1475 let (seq, rest) = be64(body)?;
1476 if !rest.is_empty() {
1477 return None;
1478 }
1479 ParsedKey::Relay(seq)
1480 }
1481 b"os" if body.is_empty() => ParsedKey::OutboxSequence,
1482 b"oc" if body.is_empty() => ParsedKey::OutcomeBacklog,
1483 b"p" => ParsedKey::Replay(hash(body)?),
1484 b"px" => {
1485 let (expires_at_ms, rest) = be64(body)?;
1486 ParsedKey::ReplayExpiry {
1487 expires_at_ms,
1488 scope: hash(rest)?,
1489 }
1490 }
1491 b"q" => ParsedKey::Quota(text(body)?),
1492 b"qx" => {
1493 let (window_start_ms, rest) = be64(body)?;
1494 ParsedKey::QuotaWindow {
1495 window_start_ms,
1496 scope: text(rest)?,
1497 }
1498 }
1499 b"qs" | b"qv" | b"qt" | b"qc" => parse_namespace_quota(tag, body)?,
1500 b"w" => {
1501 let (due_at_ms, kind, _, _, reference) = timer_parts(body)?;
1502 ParsedKey::Timer {
1503 due_at_ms,
1504 kind,
1505 reference: Bytes::copy_from_slice(reference),
1506 }
1507 }
1508 b"h" => parse_holder(body)?,
1509 b"g" => {
1510 let (object, hold_id) = body.split_first_chunk::<32>()?;
1511 ParsedKey::Hold {
1512 object: *object,
1513 hold_id: hash(hold_id)?,
1514 }
1515 }
1516 b"gp" => {
1517 let (object, hold_id) = body.split_first_chunk::<32>()?;
1518 ParsedKey::PendingHolder {
1519 object: *object,
1520 hold_id: hash(hold_id)?,
1521 }
1522 }
1523 b"ct" => {
1524 let (object, intent) = body.split_first_chunk::<32>()?;
1525 ParsedKey::ContentTakedown {
1526 object: *object,
1527 intent: hash(intent)?,
1528 }
1529 }
1530 b"b" => ParsedKey::Block(hash(body)?),
1531 b"c" => ParsedKey::ObjectState(hash(body)?),
1532 _ => return None,
1533 })
1534}
1535
1536#[must_use]
1538pub fn quota_for_window(index: &Key) -> Option<Key> {
1539 match parse(index)? {
1540 ParsedKey::QuotaWindow { scope, .. } => Some(key(TAG_QUOTA, &[scope.as_bytes()])),
1541 _ => None,
1542 }
1543}
1544
1545#[cfg(test)]
1546mod tests {
1547 use super::*;
1548 use proptest::prelude::*;
1549
1550 fn repo(name: &str) -> RepoName {
1551 RepoName::new(name).unwrap()
1552 }
1553
1554 fn scope() -> QuotaScope {
1555 QuotaScope::for_signer(&NamespaceKey::deployment_default(), &[0xab; 32])
1556 }
1557
1558 #[test]
1559 fn timer_retry_metadata_preserves_opaque_reference_and_original_due() {
1560 let reference = b"\0\xffarbitrary\0reference";
1561 let initial = timer(17, 255, reference);
1562 assert_eq!(timer_retry_state(&initial), Some((17, 0)));
1563 let retried = timer_retry(600_017, 255, reference, 17, MAX_TIMER_RETRY_ATTEMPT);
1564 assert_eq!(
1565 timer_retry_state(&retried),
1566 Some((17, MAX_TIMER_RETRY_ATTEMPT))
1567 );
1568 assert_eq!(
1569 parse(&retried),
1570 Some(ParsedKey::Timer {
1571 due_at_ms: 600_017,
1572 kind: 255,
1573 reference: Bytes::copy_from_slice(reference),
1574 })
1575 );
1576 assert!(initial < retried, "physical due remains the leading index");
1577 }
1578
1579 #[test]
1580 fn timer_retry_parser_refuses_invalid_or_truncated_metadata() {
1581 for invalid in [
1582 timer_retry(10, 12, b"work", 10, MAX_TIMER_RETRY_ATTEMPT + 1),
1583 timer_retry(10, 12, b"work", 11, 1),
1584 timer_retry(10, 12, b"work", 9, 0),
1585 Key::new(b"w\0\0\0\0\0\0\0\0\x0a\x0c\0".to_vec()),
1586 Key::new(b"w\0\0\0\0\0\0\0\0\x0a\x0cwork".to_vec()),
1587 ] {
1588 assert_eq!(parse(&invalid), None);
1589 assert_eq!(timer_retry_state(&invalid), None);
1590 }
1591 assert_eq!(timer_retry_state(&grant_epoch()), None);
1592 }
1593
1594 fn all_tags() -> Vec<&'static str> {
1595 let mut tags = vec![
1596 TAG_SHARDING_MARKER,
1597 TAG_ADDRESSING_MARKER,
1598 TAG_INSPECTION_MARKER,
1599 TAG_INSPECTION_FLAG,
1600 TAG_INSPECTION_VERSION,
1601 TAG_INSPECTION_HOLD,
1602 TAG_INSPECTION_HOLD_INDEX,
1603 TAG_INSPECTION_HOLD_MANIFEST,
1604 TAG_LAYOUT_VERSION,
1605 TAG_REF,
1606 TAG_REF_INDEX,
1607 TAG_REPLAY,
1608 TAG_REPLAY_EXPIRY,
1609 TAG_QUOTA,
1610 TAG_QUOTA_WINDOW,
1611 TAG_QUOTA_SHARD,
1612 TAG_QUOTA_VIEW,
1613 TAG_QUOTA_CONTRIBUTION,
1614 TAG_QUOTA_TOTAL,
1615 TAG_GRANT_EPOCH,
1616 TAG_EPOCH_LEASE,
1617 TAG_LEASED_SHARD,
1618 TAG_LEASE_RECOVERY,
1619 TAG_FENCE_CURSOR,
1620 TAG_BACKUP_STATE,
1621 TAG_TIMER,
1622 TAG_HOLDER,
1623 TAG_HOLD,
1624 TAG_BLOCK,
1625 TAG_OBJECT_STATE,
1626 TAG_NAMESPACE_RECORD,
1627 TAG_REPO_REGISTRY,
1628 TAG_REPO_LIST,
1629 TAG_REPO_VISIBILITY,
1630 TAG_REPO_KNOWN,
1631 TAG_RELAY_HIGH_WATER,
1632 TAG_RELAY_SCAN,
1633 TAG_TICKET,
1634 TAG_TICKET_INDEX,
1635 TAG_TICKETS_PER_REF,
1636 TAG_TICKETS_PER_SIGNER,
1637 TAG_MEMBERSHIP,
1638 TAG_VERIFICATION,
1639 TAG_VERIFY_CURSOR,
1640 TAG_OBJECT_INDEX,
1641 TAG_RESERVATION,
1642 TAG_OUTCOME_PENDING,
1643 TAG_RELAY,
1644 TAG_OUTBOX_SEQUENCE,
1645 TAG_OUTCOME_BACKLOG,
1646 TAG_CACHE_PURGE,
1647 TAG_CACHE_PURGE_GENERATION,
1648 TAG_PUBLICATION,
1649 TAG_ADVANCE,
1650 TAG_PUBLISHED_REF,
1651 TAG_PUBLISHED_INDEX,
1652 TAG_PUBLISHED_MEMBER,
1653 ];
1654 tags.extend_from_slice(RESERVED_TAGS);
1655 tags
1656 }
1657
1658 #[test]
1659 fn inspection_key_goldens_and_strict_parsing() {
1660 let repository = repo("room");
1661 let content = [0xff; 32];
1662 let advance_id = [0x02; 32];
1663 let cases = [
1664 (
1665 inspection_marker(),
1666 b"im\0".to_vec(),
1667 ParsedKey::InspectionMarker,
1668 ),
1669 (
1670 inspection_version(&repository),
1671 b"iv\0room".to_vec(),
1672 ParsedKey::InspectionVersion(repository.clone()),
1673 ),
1674 (
1675 inspection_flag(&repository, &content),
1676 [&b"if\0room\0"[..], &content].concat(),
1677 ParsedKey::InspectionFlag {
1678 repo: repository.clone(),
1679 id: content,
1680 },
1681 ),
1682 (
1683 inspection_hold(&repository, &content, &advance_id),
1684 [&b"ih\0room\0"[..], &content, &advance_id].concat(),
1685 ParsedKey::InspectionHold {
1686 repo: repository.clone(),
1687 content,
1688 advance: advance_id,
1689 },
1690 ),
1691 (
1692 inspection_hold_index(&repository, &advance_id),
1693 [&b"ia\0room\0"[..], &advance_id].concat(),
1694 ParsedKey::InspectionHoldIndex {
1695 repo: repository.clone(),
1696 advance: advance_id,
1697 },
1698 ),
1699 (
1700 inspection_hold_manifest(&repository, &advance_id),
1701 [&b"ir\0room\0"[..], &advance_id].concat(),
1702 ParsedKey::InspectionHoldManifest {
1703 repo: repository.clone(),
1704 advance: advance_id,
1705 },
1706 ),
1707 ];
1708 for (key, golden, parsed) in cases {
1709 assert_eq!(key.as_bytes(), golden);
1710 assert_eq!(parse(&key), Some(parsed));
1711 let mut extra = golden.clone();
1712 extra.push(0);
1713 if key != inspection_version(&repository) {
1714 assert_eq!(parse(&Key::new(extra)), None);
1715 }
1716 let mut short = golden;
1717 short.pop();
1718 if key != inspection_version(&repository) {
1719 assert_eq!(parse(&Key::new(short)), None);
1720 }
1721 }
1722
1723 let (start, end) = inspection_hold_range(&repository, &content);
1724 assert!(start <= inspection_hold(&repository, &content, &advance_id));
1725 assert!(inspection_hold(&repository, &content, &[0xff; 32]) < end);
1726 assert!(inspection_hold(&repo("other"), &content, &advance_id) < start);
1727 assert!(inspection_hold(&repository, &[0xfe; 32], &advance_id) < start);
1728 }
1729
1730 #[test]
1731 fn publication_key_goldens_and_strict_parsing() {
1732 let repository = repo("a");
1733 let name = "refs/heads/main";
1734 let cases = [
1735 (
1736 publication(&repository, name),
1737 b"pp\0a\0refs/heads/main".to_vec(),
1738 ParsedKey::Publication {
1739 repo: repository.clone(),
1740 name: name.into(),
1741 },
1742 ),
1743 (
1744 advance(&repository, name, 0x0102_0304_0506_0708),
1745 [
1746 b"av\0a\0refs/heads/main\0".as_slice(),
1747 &[1, 2, 3, 4, 5, 6, 7, 8],
1748 ]
1749 .concat(),
1750 ParsedKey::Advance {
1751 repo: repository.clone(),
1752 name: name.into(),
1753 sequence: 0x0102_0304_0506_0708,
1754 },
1755 ),
1756 (
1757 published_ref(&repository, name),
1758 b"pr\0a\0refs/heads/main".to_vec(),
1759 ParsedKey::PublishedRef {
1760 repo: repository.clone(),
1761 name: name.into(),
1762 },
1763 ),
1764 (
1765 published_index(&repository, name),
1766 b"py\0a\0refs/heads/main".to_vec(),
1767 ParsedKey::PublishedIndex {
1768 repo: repository.clone(),
1769 name: name.into(),
1770 },
1771 ),
1772 (
1773 published_member(&repository, &[0x11; 32]),
1774 [b"pm\0a\0".as_slice(), &[0x11; 32]].concat(),
1775 ParsedKey::PublishedMember {
1776 repo: repository.clone(),
1777 pack_id: [0x11; 32],
1778 },
1779 ),
1780 ];
1781 for (key, golden, parsed) in cases {
1782 assert_eq!(key.as_bytes(), golden);
1783 assert_eq!(parse(&key), Some(parsed));
1784 }
1785 for bad in [
1786 advance(&repository, name, 0).into_bytes().to_vec(),
1787 [b"av\0a\0refs/heads/main\0".as_slice(), &[1; 7]].concat(),
1788 [b"av\0a\0refs/heads/main\0".as_slice(), &[1; 9]].concat(),
1789 b"pr\0a\0not-a-ref".to_vec(),
1790 [b"pm\0a\0".as_slice(), &[1; 31]].concat(),
1791 [b"pm\0a\0".as_slice(), &[1; 33]].concat(),
1792 ] {
1793 assert_eq!(parse(&Key::new(bad)), None);
1794 }
1795 }
1796
1797 #[test]
1798 #[allow(clippy::too_many_lines)] fn layouts_golden_bytes() {
1800 let s = [0x11; 32];
1801 let q = format!("root\n{}", "ab".repeat(32));
1802 let cases: Vec<(Key, Vec<u8>)> = vec![
1803 (sharding_marker(), b"sm\0".to_vec()),
1804 (addressing_marker(), b"am\0".to_vec()),
1805 (inspection_marker(), b"im\0".to_vec()),
1806 (layout_version(), b"v\0".to_vec()),
1807 (relay_scan(), b"rs\0".to_vec()),
1808 (
1809 relay_high_water(&Partition::ContentShard(7)).unwrap(),
1810 b"rh\0s7\0".to_vec(),
1811 ),
1812 (namespace_record(), b"nr\0".to_vec()),
1813 (repo_record(&repo("room-a")), b"rr\0room-a".to_vec()),
1814 (repo_visibility(&repo("room-a")), b"rv\0room-a".to_vec()),
1815 (repo_known(&repo("room-a")), b"rk\0room-a".to_vec()),
1816 (
1817 verification(&repo("room-a"), &s),
1818 [&b"vs\0room-a\0"[..], &[0x11; 32]].concat(),
1819 ),
1820 (
1821 ref_key(&repo("room-a"), "refs/heads/main"),
1822 b"r\0room-a\0refs/heads/main".to_vec(),
1823 ),
1824 (
1825 ref_index_key(&repo("room-a"), "refs/heads/main"),
1826 b"x\0room-a\0refs/heads/main".to_vec(),
1827 ),
1828 (replay(&s), [&b"p\0"[..], &[0x11; 32]].concat()),
1829 (
1830 replay_expiry(0x0102_0304_0506_0708, &s),
1831 [&b"px\0"[..], &[1, 2, 3, 4, 5, 6, 7, 8], &[0x11; 32]].concat(),
1832 ),
1833 (quota(&scope()), [b"q\0", q.as_bytes()].concat()),
1834 (
1835 quota_window(256, &scope()),
1836 [&b"qx\0"[..], &[0, 0, 0, 0, 0, 0, 1, 0], q.as_bytes()].concat(),
1837 ),
1838 (grant_epoch(), b"e\0".to_vec()),
1839 (epoch_lease(), b"el\0".to_vec()),
1840 (lease_recovery(), b"lr\0".to_vec()),
1841 (lease_reconcile(), b"lrc\0".to_vec()),
1842 (backup_state(), b"bk\0".to_vec()),
1843 (
1844 leased_shard(&repo("a"), "refs/heads/main"),
1845 b"ls\0a\0refs/heads/main".to_vec(),
1846 ),
1847 (
1848 timer(1, 7, b"refs/heads/x"),
1849 [
1850 &b"w\0"[..],
1851 &[0, 0, 0, 0, 0, 0, 0, 1, 7, 0],
1852 &[0, 0, 0, 0, 0, 0, 0, 1],
1853 b"refs/heads/x",
1854 ]
1855 .concat(),
1856 ),
1857 (
1858 holder(&s, &NamespaceKey::deployment_default(), &repo("a")).unwrap(),
1859 [&b"h\0"[..], &[0x11; 32], b"root\0a"].concat(),
1860 ),
1861 (
1862 hold(&s, &[0x22; 32]),
1863 [&b"g\0"[..], &[0x11; 32], &[0x22; 32]].concat(),
1864 ),
1865 (block(&s), [&b"b\0"[..], &[0x11; 32]].concat()),
1866 (object_state(&s), [&b"c\0"[..], &[0x11; 32]].concat()),
1867 (
1868 object_index(&repo("a"), &s, &[0x22; 32]),
1869 [&b"i\0a\0"[..], &[0x11; 32], &[0x22; 32]].concat(),
1870 ),
1871 ];
1872 for (key, golden) in cases {
1873 assert_eq!(key.as_bytes(), golden.as_slice());
1874 }
1875 let (a, s2) = (repo("a"), [0x22; 32]);
1876 for (key, golden, sub, id) in [
1877 (
1878 verify_job(&a, &s),
1879 [&b"vc\0a\0"[..], &s, &[0]].concat(),
1880 VC_JOB,
1881 None,
1882 ),
1883 (
1884 verify_row(&a, &s, VC_FRAME, Some(&s2)),
1885 [&b"vc\0a\0"[..], &s, &[1], &s2].concat(),
1886 VC_FRAME,
1887 Some(s2),
1888 ),
1889 (
1890 verify_row(&a, &s, VC_CANDIDATE, Some(&s2)),
1891 [&b"vc\0a\0"[..], &s, &[4], &s2].concat(),
1892 VC_CANDIDATE,
1893 Some(s2),
1894 ),
1895 (
1896 verify_row(&a, &s, VC_DEPENDENCY, Some(&s2)),
1897 [&b"vc\0a\0"[..], &s, &[6], &s2].concat(),
1898 VC_DEPENDENCY,
1899 Some(s2),
1900 ),
1901 ] {
1902 assert_eq!(key.as_bytes(), golden.as_slice());
1903 assert_eq!(
1904 parse(&key),
1905 Some(ParsedKey::VerifyCursor {
1906 repo: a.clone(),
1907 pack_id: s,
1908 sub,
1909 id,
1910 })
1911 );
1912 }
1913 let (start, end) = verify_range(&a, &s, None);
1914 for sub in [VC_JOB, VC_FRAME, VC_BASE, VC_DEPENDENCY] {
1915 let row = verify_row(&a, &s, sub, (sub != VC_JOB).then_some(&s2));
1916 assert!(start <= row && row < end);
1917 }
1918 let (start, end) = verify_range(&a, &s, Some(VC_CHILD));
1919 assert!(
1920 !(start <= verify_row(&a, &s, VC_FRAME, Some(&s2))
1921 && verify_row(&a, &s, VC_FRAME, Some(&s2)) < end)
1922 );
1923 assert!(start <= verify_row(&a, &s, VC_CHILD, Some(&s2)));
1924 let other = verify_job(&a, &s2);
1925 assert!(!(verify_range(&a, &s, None).0 <= other && other < verify_range(&a, &s, None).1));
1926 for bad in [
1927 [&b"vc\0a\0"[..], &s, &[0], &s2].concat(),
1928 [&b"vc\0a\0"[..], &s, &[1]].concat(),
1929 [&b"vc\0a\0"[..], &s, &[7], &s2].concat(),
1930 ] {
1931 assert_eq!(parse(&Key::new(bad)), None);
1932 }
1933 assert_eq!(LAYOUT_VERSION, 1);
1934 assert!(!RESERVED_TAGS.contains(&TAG_VERIFY_CURSOR));
1935 assert!(!RESERVED_TAGS.contains(&TAG_OBJECT_INDEX));
1936 assert!(!RESERVED_TAGS.contains(&TAG_VERIFICATION));
1937 let state = verification(&repo("a"), &s);
1938 assert_eq!(
1939 parse(&state),
1940 Some(ParsedKey::Verification {
1941 repo: repo("a"),
1942 pack_id: s,
1943 })
1944 );
1945 let index = object_index(&repo("a"), &s, &[0x22; 32]);
1946 assert_eq!(
1947 parse(&index),
1948 Some(ParsedKey::ObjectIndex {
1949 repo: repo("a"),
1950 object: s,
1951 pack_id: [0x22; 32],
1952 })
1953 );
1954 let (start, end) = object_index_range(&repo("a"), &s);
1955 assert!(start <= index && index < end);
1956 assert_eq!(
1957 parse(&Key::new([&b"i\0a\0"[..], &[0x11; 63]].concat())),
1958 None
1959 );
1960 for (tag, start, end) in [
1962 (TAG_REPO_REGISTRY, &b"rr\0"[..], &b"rr\x01"[..]),
1963 (TAG_NAMESPACE_LIST, b"nl\0", b"nl\x01"),
1964 ] {
1965 let (s, e) = class_range(tag);
1966 assert_eq!((s.as_bytes(), e.as_bytes()), (start, end));
1967 }
1968 }
1969
1970 #[test]
1971 fn ticket_outbox_layouts_golden_and_roundtrip() {
1972 let r = repo("a");
1973 let name = "refs/heads/main";
1974 let pack = [0x11; 32];
1975 let signer = [0x22; 32];
1976 let seq = 0x0102_0304_0506_0708;
1977 let cases = vec![
1978 (
1979 ticket(&pack),
1980 [&b"t\0"[..], &pack].concat(),
1981 ParsedKey::Ticket(pack),
1982 ),
1983 (
1984 ticket_index(&r, name, &pack, &signer).unwrap(),
1985 [&b"ti\0a\0refs/heads/main\0"[..], &pack, &signer].concat(),
1986 ParsedKey::TicketIndex {
1987 repo: r.clone(),
1988 name: name.into(),
1989 pack_id: pack,
1990 signer,
1991 },
1992 ),
1993 (
1994 tickets_per_ref(&r, name).unwrap(),
1995 b"tc\0a\0refs/heads/main".to_vec(),
1996 ParsedKey::TicketsPerRef {
1997 repo: r.clone(),
1998 name: name.into(),
1999 },
2000 ),
2001 (
2002 tickets_per_signer(&r, name, &signer).unwrap(),
2003 [&b"tu\0a\0refs/heads/main\0"[..], &signer].concat(),
2004 ParsedKey::TicketsPerSigner {
2005 repo: r.clone(),
2006 name: name.into(),
2007 signer,
2008 },
2009 ),
2010 (
2011 membership(&r, &pack),
2012 [&b"m\0a\0"[..], &pack].concat(),
2013 ParsedKey::Membership {
2014 repo: r,
2015 pack_id: pack,
2016 },
2017 ),
2018 (
2019 reservation("R-1:ok").unwrap(),
2020 b"o\0R-1:ok".to_vec(),
2021 ParsedKey::Reservation("R-1:ok".into()),
2022 ),
2023 (
2024 outcome_pending(seq, "R-1:ok").unwrap(),
2025 [&b"oq\0"[..], &[1, 2, 3, 4, 5, 6, 7, 8], b"R-1:ok"].concat(),
2026 ParsedKey::OutcomePending {
2027 seq,
2028 reservation_id: "R-1:ok".into(),
2029 },
2030 ),
2031 (
2032 relay(seq),
2033 [&b"or\0"[..], &[1, 2, 3, 4, 5, 6, 7, 8]].concat(),
2034 ParsedKey::Relay(seq),
2035 ),
2036 (relay_scan(), b"rs\0".to_vec(), ParsedKey::RelayScan),
2037 (
2038 outbox_sequence(),
2039 b"os\0".to_vec(),
2040 ParsedKey::OutboxSequence,
2041 ),
2042 (
2043 outcome_backlog(),
2044 b"oc\0".to_vec(),
2045 ParsedKey::OutcomeBacklog,
2046 ),
2047 (
2048 timer(seq, 2, &pack),
2049 [&b"w\0\x01\x02\x03\x04\x05\x06\x07\x08\x02\0\x01\x02\x03\x04\x05\x06\x07\x08"[..], &pack].concat(),
2050 ParsedKey::Timer {
2051 due_at_ms: seq,
2052 kind: 2,
2053 reference: Bytes::copy_from_slice(&pack),
2054 },
2055 ),
2056 ];
2057 for (key, golden, parsed) in cases {
2058 assert_eq!(key.as_bytes(), golden);
2059 assert_eq!(parse(&key), Some(parsed));
2060 }
2061 assert_eq!(parse(&Key::new(b"rs\0extra".to_vec())), None);
2062 let (start, end) = class_range(TAG_RELAY_SCAN);
2063 assert_eq!(
2064 (start.as_bytes(), end.as_bytes()),
2065 (&b"rs\0"[..], &b"rs\x01"[..])
2066 );
2067 }
2068
2069 #[test]
2070 fn ticket_outbox_key_validation_and_maximum() {
2071 let r = repo(&"a".repeat(MAX_REPO_NAME_BYTES));
2072 let name = format!(
2073 "refs/heads/{}",
2074 "b".repeat(MAX_REF_NAME_BYTES - "refs/heads/".len())
2075 );
2076 let longest = ticket_index(&r, &name, &[0; 32], &[0; 32]).unwrap();
2077 assert_eq!(
2078 longest.as_bytes().len(),
2079 3 + MAX_REPO_NAME_BYTES + 1 + MAX_REF_NAME_BYTES + 1 + 64
2080 );
2081 assert!(longest.as_bytes().len() <= MAX_KEY_BYTES);
2082 assert!(parse(&longest).is_some());
2083 for rid in ["", "bad/rid", "space id", "nonascii-é", &"a".repeat(129)] {
2084 assert!(!validate_reservation_id(rid));
2085 assert!(reservation(rid).is_err());
2086 assert!(outcome_pending(1, rid).is_err());
2087 }
2088 assert!(reservation(&"a".repeat(128)).is_ok());
2089 assert!(outcome_pending(0, "ok").is_err());
2090 for name in ["", "bad", "refs/heads/a\0b", &format!("{name}x")] {
2091 assert!(ticket_index(&r, name, &[0; 32], &[0; 32]).is_err());
2092 assert!(tickets_per_ref(&r, name).is_err());
2093 assert!(tickets_per_signer(&r, name, &[0; 32]).is_err());
2094 }
2095 for bad in [
2096 &b"t\0short"[..],
2097 b"ti\0a\0refs/heads/a\0short",
2098 b"tc\0a\0bad",
2099 b"tu\0a\0refs/heads/a\0short",
2100 b"m\0a\0short",
2101 b"o\0bad/id",
2102 b"oq\0short",
2103 b"or\0short",
2104 b"os\0extra",
2105 b"oc\0extra",
2106 ] {
2107 assert_eq!(parse(&Key::new(bad.to_vec())), None);
2108 }
2109 assert_eq!(parse(&Key::new(vec![b'x'; MAX_KEY_BYTES + 1])), None);
2110 }
2111
2112 #[test]
2113 fn backup_state_parse_roundtrip_and_malformed_suffix() {
2114 assert_eq!(parse(&backup_state()), Some(ParsedKey::BackupState));
2115 assert_eq!(parse(&Key::new(b"bk\0x".to_vec())), None);
2116 }
2117
2118 #[test]
2119 fn class_scans_never_overlap() {
2120 let tags = all_tags();
2121 for (i, a) in tags.iter().enumerate() {
2122 assert!(!tags[i + 1..].contains(a), "duplicate tag {a}");
2123 let (start, end) = class_range(a);
2124 for b in &tags {
2125 for tail in [&b""[..], b"\0", b"\xff\xff", b"x\0y"] {
2126 let k = Key::new([b.as_bytes(), b"\0", tail].concat());
2127 assert_eq!(start <= k && k < end, a == b, "{a} range vs {b} key");
2128 }
2129 }
2130 }
2131 }
2132
2133 #[test]
2134 fn ref_index_prefix_is_exact() {
2135 let first = repo("one");
2136 let (start, end) = ref_index_prefix_range(&first, "refs/heads/feat/");
2137 assert!(
2138 start <= ref_index_key(&first, "refs/heads/feat/topic")
2139 && ref_index_key(&first, "refs/heads/feat/topic") < end
2140 );
2141 assert!(
2142 !(start <= ref_index_key(&first, "refs/heads/featx")
2143 && ref_index_key(&first, "refs/heads/featx") < end)
2144 );
2145 assert!(
2146 !(start <= ref_index_key(&repo("two"), "refs/heads/feat/topic")
2147 && ref_index_key(&repo("two"), "refs/heads/feat/topic") < end)
2148 );
2149 }
2150
2151 proptest! {
2152 #[test]
2153 fn ref_index_prefix_bounds_and_repo_isolation(
2154 prefix in "[a-z/.-]{0,6}",
2155 name in "[a-z/.-]{0,10}",
2156 ) {
2157 let r = repo("repo");
2158 let (start, end) = ref_index_prefix_range(&r, &prefix);
2159 let key = ref_index_key(&r, &name);
2160 prop_assert_eq!(start <= key && key < end, name.starts_with(&prefix));
2161 let foreign = ref_index_key(&repo("repo2"), &name);
2162 prop_assert!(!(start <= foreign && foreign < end));
2163 prop_assert_eq!(parse(&key), Some(ParsedKey::RefIndexEntry { repo: r, name }));
2164 }
2165
2166 #[test]
2167 fn ref_prefix_scan_bounds_cover_exactly_the_prefix(
2168 prefix in "[a-z/.-]{0,6}",
2169 name in "[a-z/.-]{0,10}",
2170 other in "[a-z]{1,3}",
2171 ) {
2172 let r = repo("repo");
2173 let (start, end) = ref_prefix_range(&r, &prefix);
2174 let k = ref_key(&r, &name);
2175 prop_assert_eq!(start <= k && k < end, name.starts_with(&prefix));
2176 let foreign = ref_key(&repo(&format!("repo{other}")), &name);
2178 prop_assert!(!(start <= foreign && foreign < end));
2179 prop_assert!(!(start <= replay(&[0; 32]) && replay(&[0; 32]) < end));
2180 }
2181
2182 #[test]
2183 fn be64_orders_numerically(a: u64, b: u64) {
2184 let s = [0xff; 32];
2185 prop_assert_eq!(replay_expiry(a, &s).cmp(&replay_expiry(b, &s)), a.cmp(&b));
2186 let (start, end) = replay_expiry_before(b);
2187 let k = replay_expiry(a, &s);
2188 prop_assert_eq!(start <= k && k < end, a < b);
2189 let (start, end) = quota_window_before(b);
2190 let k = quota_window(a, &scope());
2191 prop_assert_eq!(start <= k && k < end, a < b);
2192 }
2193 }
2194
2195 #[test]
2196 #[allow(clippy::too_many_lines)] fn parse_roundtrip_every_class() {
2198 let s = [0x22; 32];
2199 let q = scope();
2200 let cases = vec![
2201 (layout_version(), ParsedKey::LayoutVersion),
2202 (
2203 ref_key(&repo("a"), "refs/tags/v1"),
2204 ParsedKey::Ref {
2205 repo: repo("a"),
2206 name: "refs/tags/v1".into(),
2207 },
2208 ),
2209 (replay(&s), ParsedKey::Replay(s)),
2210 (
2211 replay_expiry(9, &s),
2212 ParsedKey::ReplayExpiry {
2213 expires_at_ms: 9,
2214 scope: s,
2215 },
2216 ),
2217 (quota(&q), ParsedKey::Quota(q.as_str().into())),
2218 (
2219 quota_window(5, &q),
2220 ParsedKey::QuotaWindow {
2221 window_start_ms: 5,
2222 scope: q.as_str().into(),
2223 },
2224 ),
2225 (grant_epoch(), ParsedKey::GrantEpoch),
2226 (epoch_lease(), ParsedKey::EpochLease),
2227 (lease_recovery(), ParsedKey::LeaseRecovery),
2228 (revoke_cursor(false), ParsedKey::FenceCursor(0)),
2229 (revoke_cursor(true), ParsedKey::FenceCursor(1)),
2230 (lease_reconcile(), ParsedKey::LeaseReconcile),
2231 (
2232 leased_shard(&repo("a"), "refs/heads/main"),
2233 ParsedKey::LeasedShard {
2234 repo: repo("a"),
2235 shard_ref: "refs/heads/main".into(),
2236 },
2237 ),
2238 (sharding_marker(), ParsedKey::ShardingMarker),
2239 (addressing_marker(), ParsedKey::AddressingMarker),
2240 (namespace_record(), ParsedKey::NamespaceRecord),
2241 (repo_record(&repo("a")), ParsedKey::RepoRecord(repo("a"))),
2242 (repo_known(&repo("a")), ParsedKey::RepoKnown(repo("a"))),
2243 (
2244 timer(3, 2, b"r"),
2245 ParsedKey::Timer {
2246 due_at_ms: 3,
2247 kind: 2,
2248 reference: Bytes::from_static(b"r"),
2249 },
2250 ),
2251 (
2252 holder(&s, &NamespaceKey::deployment_default(), &repo("a")).unwrap(),
2253 ParsedKey::Holder {
2254 object: s,
2255 ns: NamespaceKey::deployment_default(),
2256 repo: repo("a"),
2257 },
2258 ),
2259 (
2260 hold(&s, &[3; 32]),
2261 ParsedKey::Hold {
2262 object: s,
2263 hold_id: [3; 32],
2264 },
2265 ),
2266 (block(&s), ParsedKey::Block(s)),
2267 (object_state(&s), ParsedKey::ObjectState(s)),
2268 ];
2269 for (key, parsed) in cases {
2270 assert_eq!(parse(&key), Some(parsed));
2271 }
2272 assert_eq!(quota_for_window("a_window(5, &q)), Some(quota(&q)));
2273 assert_eq!(quota_for_window("a(&q)), None);
2274 for bad in [
2275 &b"p\0short"[..],
2276 b"v\0x",
2277 b"px\0\0",
2278 b"no-terminator",
2279 b"c\0short",
2280 b"g\0short",
2281 b"nr\0extra",
2282 b"rr\0",
2283 b"rk\0bad name",
2284 ] {
2285 assert_eq!(parse(&Key::new(bad.to_vec())), None);
2286 }
2287 let other = [0x23; 32];
2289 let (start, end) = holders_of(&s);
2290 let own = holder(&s, &NamespaceKey::deployment_default(), &repo("z")).unwrap();
2291 let foreign = holder(&other, &NamespaceKey::deployment_default(), &repo("a")).unwrap();
2292 assert!(start <= own && own < end && !(start <= foreign && foreign < end));
2293 let (start, end) = holds_of(&s);
2294 assert!(start <= hold(&s, &[0xff; 32]) && hold(&s, &[0xff; 32]) < end);
2295 assert!(!(start <= hold(&other, &[0; 32]) && hold(&other, &[0; 32]) < end));
2296 let nul = NamespaceKey::from_stored("a\0b".into());
2297 assert!(matches!(
2298 holder(&s, &nul, &repo("a")),
2299 Err(StoreError::Invalid(_))
2300 ));
2301 }
2302
2303 #[test]
2304 fn object_index_parse_roundtrip() {
2305 let object = [0x22; 32];
2306 let pack_id = [3; 32];
2307 let key = object_index(&repo("a"), &object, &pack_id);
2308 assert_eq!(
2309 parse(&key),
2310 Some(ParsedKey::ObjectIndex {
2311 repo: repo("a"),
2312 object,
2313 pack_id,
2314 })
2315 );
2316 }
2317 #[test]
2318 fn repo_listing_subkeys_are_sorted_and_strict() {
2319 let repo = RepoName::new("repo-a").unwrap();
2320 for explicit_public in [false, true] {
2321 let key = repo_listing(&repo, explicit_public);
2322 assert!(
2323 key.as_bytes()
2324 .starts_with(repo_listing_prefix(explicit_public).as_bytes())
2325 );
2326 assert_eq!(
2327 parse(&key),
2328 Some(ParsedKey::RepoListing {
2329 explicit_public,
2330 repo: repo.clone()
2331 })
2332 );
2333 assert!(key < repo_listing(&RepoName::new("repo-b").unwrap(), explicit_public));
2334 }
2335 for bytes in [b"rl\0x\0repo-a".as_slice(), b"rl\0p\0", b"rl\0d\0bad\0name"] {
2336 assert_eq!(parse(&Key::new(bytes.to_vec())), None);
2337 }
2338 }
2339
2340 #[test]
2341 fn repo_visibility_key_roundtrips() {
2342 let key = repo_visibility(&repo("room-a"));
2343 assert_eq!(key.as_bytes(), b"rv\0room-a");
2344 assert_eq!(parse(&key), Some(ParsedKey::RepoVisibility(repo("room-a"))));
2345 assert_eq!(parse(&Key::new(b"rv\0"[..].to_vec())), None);
2346 }
2347
2348 #[test]
2349 fn lease_keys_reject_malformed_payloads() {
2350 for bad in [
2351 &b"el\0extra"[..],
2352 b"lr\0extra",
2353 b"ls\0no-separator",
2354 b"ls\0\0refs/heads/main",
2355 ] {
2356 assert_eq!(parse(&Key::new(bad.to_vec())), None);
2357 }
2358 }
2359
2360 #[test]
2361 fn namespace_quota_key_goldens_and_ranges() {
2362 let window: u64 = 0x0102_0304_0506_0708;
2363 let suffix = window.to_be_bytes();
2364 assert_eq!(
2365 quota_shard(window).as_bytes(),
2366 [&b"qs\0"[..], &suffix].concat()
2367 );
2368 assert_eq!(
2369 quota_view(window).as_bytes(),
2370 [&b"qv\0"[..], &suffix].concat()
2371 );
2372 assert_eq!(
2373 quota_total(window).as_bytes(),
2374 [&b"qt\0"[..], &suffix].concat()
2375 );
2376 let source = Partition::Ref {
2377 ns: NamespaceKey::deployment_default(),
2378 repo: repo("room"),
2379 shard_ref: "refs/heads/main".into(),
2380 };
2381 let contribution = quota_contribution(window, &source).unwrap();
2382 let next = quota_contribution(window + 1, &source).unwrap();
2383 assert_eq!(
2384 contribution.as_bytes(),
2385 [&b"qc\0"[..], &suffix, &source.encode().unwrap()].concat()
2386 );
2387 assert_eq!(
2388 parse(&contribution),
2389 Some(ParsedKey::QuotaContribution { window, source })
2390 );
2391 assert_eq!(
2392 parse("a_shard(window)),
2393 Some(ParsedKey::QuotaShard(window))
2394 );
2395 assert_eq!(
2396 parse("a_view(window)),
2397 Some(ParsedKey::QuotaView(window))
2398 );
2399 assert_eq!(
2400 parse("a_total(window)),
2401 Some(ParsedKey::QuotaTotal(window))
2402 );
2403 let (start, end) = quota_namespace_window(TAG_QUOTA_CONTRIBUTION, window);
2404 assert!(start <= contribution && contribution < end);
2405 assert!(!(start <= next && next < end));
2406 assert!(matches!(
2407 quota_contribution(
2408 window,
2409 &Partition::Coordinator(NamespaceKey::deployment_default())
2410 ),
2411 Err(StoreError::Invalid(_))
2412 ));
2413 }
2414}