1use mkit_core::hash::{Hash, from_hex, to_hex, to_hex_bytes};
7use mkit_core::protocol::AdvanceOutcome;
8use serde::{Deserialize, Serialize};
9
10use super::content_index::{BlockEntry, HolderRecord, ObjectState};
11use super::error::StoreError;
12use super::index::IndexValue;
13use super::keys::validate_reservation_id;
14use super::kv::{Key, MAX_KEY_BYTES, MAX_VALUE_BYTES, Value};
15use super::partition::Partition;
16use crate::error::Code;
17use crate::quota::{NamespaceUsage, NamespaceView, QuotaState};
18use crate::refs::is_served_ref_name;
19use crate::replay::{
20 BeginUploadResult, ReplayRecord, ReplayState, StoredRejection, StoredResult, UpdateRefResult,
21};
22use crate::repo::RepoName;
23use mkit_core::repo_identity::RepositoryIdentity;
24use mkit_core::upload_parts::MIN_PART_SIZE;
25use mkit_core::write_auth::is_hex;
26
27#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
29#[serde(deny_unknown_fields)]
30pub struct NamespaceRecord {
31 pub created_at_ms: u64,
33 pub config_version: u64,
35}
36
37#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
39#[serde(deny_unknown_fields)]
40pub struct RepoRecord {
41 pub created_at_ms: u64,
43}
44
45#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
48#[serde(deny_unknown_fields)]
49pub struct RepoVisibilityV1 {
50 pub visibility: StoredVisibility,
52 pub last_created_ms: u64,
55 #[serde(deserialize_with = "present_option")]
59 pub last_statement_id: Option<String>,
60 #[serde(default, skip_serializing_if = "Option::is_none")]
65 pub changed_ms: Option<u64>,
66}
67
68impl RepoVisibilityV1 {
69 #[must_use]
74 pub fn visibility_changed_ms(&self) -> u64 {
75 self.changed_ms.unwrap_or(self.last_created_ms)
76 }
77}
78
79fn present_option<'de, D, T>(deserializer: D) -> Result<Option<T>, D::Error>
82where
83 D: serde::Deserializer<'de>,
84 T: Deserialize<'de>,
85{
86 Option::<T>::deserialize(deserializer)
87}
88
89#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
91#[serde(rename_all = "lowercase")]
92pub enum StoredVisibility {
93 Public,
95 Private,
97}
98
99#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
101#[serde(deny_unknown_fields)]
102pub struct EpochLease {
103 #[serde(default, skip_serializing_if = "Option::is_none")]
105 pub authority_ready: Option<bool>,
106 pub epoch: u64,
108 pub expires_at_ms: u64,
110 pub config_version: u64,
112 #[serde(default, skip_serializing_if = "Option::is_none")]
114 pub authority_generation: Option<u64>,
115}
116
117#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
125#[serde(deny_unknown_fields)]
126pub struct LeasedShard {
127 pub epoch: u64,
129 pub expires_at_ms: u64,
131 pub acked_epoch: u64,
133 #[serde(default, skip_serializing_if = "Option::is_none")]
135 pub authority_generation: Option<u64>,
136 #[serde(default, skip_serializing_if = "Option::is_none")]
138 pub acked_authority_generation: Option<u64>,
139 pub relay_watermark_ms: u64,
141 pub sweep_due_ms: u64,
143}
144
145#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
147#[serde(deny_unknown_fields)]
148pub struct LeaseRecovery {
149 #[serde(default, skip_serializing_if = "Option::is_none")]
151 pub authority_fence: Option<bool>,
152 #[serde(default, skip_serializing_if = "Option::is_none")]
154 pub authority_ready: Option<bool>,
155 #[serde(default, skip_serializing_if = "Option::is_none")]
157 pub activation_only: Option<bool>,
158 pub resumed_at_ms: u64,
160}
161
162impl LeaseRecovery {
163 #[must_use]
165 pub fn recovery_time(self) -> Option<u64> {
166 (self.activation_only != Some(true)).then_some(self.resumed_at_ms)
167 }
168}
169
170#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
173#[serde(deny_unknown_fields)]
174pub struct BackupStateV1 {
175 pub last_export_ms: u64,
177 pub digest: Hash,
179 pub r2_key: String,
181 pub last_upload_ms: u64,
183}
184
185#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
187#[serde(deny_unknown_fields)]
188pub struct TicketV1 {
189 #[serde(default, skip_serializing_if = "Option::is_none")]
191 pub authority_generation: Option<u64>,
192 #[serde(with = "repo_json")]
194 pub repo: RepoName,
195 pub ref_name: String,
197 #[serde(with = "hash_json")]
199 pub signer: Hash,
200 #[serde(with = "hash_json")]
202 pub pack_id: Hash,
203 pub bytes: u64,
205 pub part_size: u64,
207 pub expires_at_ms: u64,
209 pub created_at_ms: u64,
211 pub reservation_id: String,
213 #[serde(with = "optional_bytes_hex_json")]
215 pub upload_session: Option<Vec<u8>>,
216}
217
218#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
220#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
221pub enum AbortReason {
222 Unspecified,
224 RefConflict,
226 EpochMismatch,
228 PackMissing,
230 ReplayRace,
232 Internal,
234 Abandoned,
236}
237
238#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
240#[serde(rename_all = "snake_case")]
241pub enum PendingOp {
242 Write,
244 Read,
246}
247
248#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
250#[serde(deny_unknown_fields)]
251pub struct OutcomeRef {
252 pub name: String,
254 #[serde(with = "optional_hash_json")]
256 pub new: Option<Hash>,
257 pub deleted: bool,
259}
260
261#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
265#[serde(tag = "state", rename_all = "snake_case", deny_unknown_fields)]
266pub enum ReservationV1 {
267 Pending {
269 repository: String,
271 created_at_ms: u64,
273 reconcile_at_ms: u64,
275 op: PendingOp,
277 },
278 Ticketed {
280 #[serde(with = "hash_json")]
282 ticket_id: Hash,
283 },
284 Committed {
286 repository: String,
288 occurred_at_ms: u64,
290 bytes_stored: u64,
292 new_to_repo: u64,
294 new_to_store: u64,
296 refs: Vec<OutcomeRef>,
298 },
299 Aborted {
301 repository: String,
303 occurred_at_ms: u64,
305 reason: AbortReason,
307 detail: String,
309 },
310 Expired {
312 repository: String,
314 occurred_at_ms: u64,
316 },
317 ReadServed {
319 repository: String,
321 occurred_at_ms: u64,
323 #[serde(with = "hash_json")]
325 object: Hash,
326 bytes_served: u64,
328 },
329}
330
331#[derive(Debug, Clone, PartialEq, Eq)]
333pub struct RelayV1 {
334 pub at_ms: u64,
336 pub target: Partition,
338 pub puts: Vec<(Key, Value)>,
340 pub deletes: Vec<Key>,
342}
343
344pub const MAX_BLOCKED_TARGETS: usize = 32;
346
347#[derive(Debug, Clone, PartialEq, Eq)]
352pub struct RelayScanV1 {
353 pub cycle_end: u64,
355 pub cursor: u64,
357 pub blocked: Vec<Partition>,
359}
360
361#[derive(Serialize, Deserialize)]
362#[serde(deny_unknown_fields)]
363struct RelayScanDtoV1 {
364 cycle_end: u64,
365 cursor: u64,
366 blocked: Vec<String>,
367}
368
369#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
371#[serde(deny_unknown_fields)]
372pub struct Backlog {
373 pub rows: u64,
375 pub bytes: u64,
377}
378
379#[derive(Serialize, Deserialize)]
380#[serde(deny_unknown_fields)]
381struct RelayDtoV1 {
382 at_ms: u64,
383 target: String,
384 puts: Vec<(String, String)>,
385 #[serde(default, skip_serializing_if = "Vec::is_empty")]
386 deletes: Vec<String>,
387}
388
389mod hash_json {
390 use super::{Deserialize, Hash, from_hex, to_hex};
391 pub(super) fn serialize<S: serde::Serializer>(
392 hash: &Hash,
393 serializer: S,
394 ) -> Result<S::Ok, S::Error> {
395 serializer.serialize_str(&to_hex(hash))
396 }
397 pub(super) fn deserialize<'de, D: serde::Deserializer<'de>>(
398 deserializer: D,
399 ) -> Result<Hash, D::Error> {
400 let s = String::deserialize(deserializer)?;
401 from_hex(&s).map_err(serde::de::Error::custom)
402 }
403}
404
405mod optional_hash_json {
406 use super::{Deserialize, Hash, Serialize, from_hex, to_hex};
407 #[allow(clippy::ref_option)]
409 pub(super) fn serialize<S: serde::Serializer>(
410 hash: &Option<Hash>,
411 serializer: S,
412 ) -> Result<S::Ok, S::Error> {
413 hash.as_ref().map(to_hex).serialize(serializer)
414 }
415 pub(super) fn deserialize<'de, D: serde::Deserializer<'de>>(
416 deserializer: D,
417 ) -> Result<Option<Hash>, D::Error> {
418 Option::<String>::deserialize(deserializer)?
419 .as_deref()
420 .map(from_hex)
421 .transpose()
422 .map_err(serde::de::Error::custom)
423 }
424}
425
426mod optional_bytes_hex_json {
427 use super::{Deserialize, Serialize, hex_nibble, to_hex_bytes};
428
429 #[allow(clippy::ref_option)]
430 pub(super) fn serialize<S: serde::Serializer>(
431 bytes: &Option<Vec<u8>>,
432 serializer: S,
433 ) -> Result<S::Ok, S::Error> {
434 bytes
435 .as_ref()
436 .map(|value| to_hex_bytes(value))
437 .serialize(serializer)
438 }
439
440 pub(super) fn deserialize<'de, D: serde::Deserializer<'de>>(
441 deserializer: D,
442 ) -> Result<Option<Vec<u8>>, D::Error> {
443 Option::<String>::deserialize(deserializer)?
444 .map(|s| {
445 if s.len() % 2 != 0 {
446 return Err(serde::de::Error::custom("odd-length upload session hex"));
447 }
448 s.as_bytes()
449 .chunks_exact(2)
450 .map(|pair| {
451 let high = hex_nibble(pair[0]).ok_or_else(|| {
452 serde::de::Error::custom("invalid upload session hex")
453 })?;
454 let low = hex_nibble(pair[1]).ok_or_else(|| {
455 serde::de::Error::custom("invalid upload session hex")
456 })?;
457 Ok((high << 4) | low)
458 })
459 .collect()
460 })
461 .transpose()
462 }
463}
464
465fn hex_nibble(byte: u8) -> Option<u8> {
466 match byte {
467 b'0'..=b'9' => Some(byte - b'0'),
468 b'a'..=b'f' => Some(byte - b'a' + 10),
469 b'A'..=b'F' => Some(byte - b'A' + 10),
470 _ => None,
471 }
472}
473
474mod repo_json {
475 use super::{Deserialize, RepoName};
476 pub(super) fn serialize<S: serde::Serializer>(
477 repo: &RepoName,
478 serializer: S,
479 ) -> Result<S::Ok, S::Error> {
480 serializer.serialize_str(repo.as_str())
481 }
482 pub(super) fn deserialize<'de, D: serde::Deserializer<'de>>(
483 deserializer: D,
484 ) -> Result<RepoName, D::Error> {
485 RepoName::new(String::deserialize(deserializer)?).map_err(serde::de::Error::custom)
486 }
487}
488
489pub const CODEC_V1: u8 = 0x01;
491
492#[derive(Serialize, Deserialize)]
493#[serde(deny_unknown_fields)]
494struct RecordV1 {
495 fingerprint: String,
496 expires_at_ms: i64,
497 state: StateV1,
498}
499
500#[derive(Serialize, Deserialize)]
501#[serde(tag = "state", rename_all = "snake_case")]
502enum StateV1 {
503 InFlight { resumable: bool },
504 Committed { result: ResultV1 },
505}
506
507#[derive(Serialize, Deserialize)]
508#[serde(tag = "kind", rename_all = "snake_case")]
509enum ResultV1 {
510 UpdateRefCommitted,
511 UpdateRefConflict {
512 current: Option<String>,
513 },
514 AdvanceCommitted,
515 AdvanceHeadConflict,
516 AdvancePackmapConflict,
517 UploadPack,
518 RepoVisibility,
519 BeginUploadAlreadyPresent,
520 BeginUploadTicket {
521 id: String,
522 part_size: u64,
523 expires_at_ms: u64,
524 token_hex: String,
525 },
526 Rejected {
527 code: String,
528 message: String,
529 },
530}
531
532#[derive(Serialize, Deserialize)]
533#[serde(deny_unknown_fields)]
534struct QuotaV1 {
535 window_start: i64,
536 ops: u32,
537 bytes: u64,
538}
539
540#[derive(Serialize, Deserialize)]
541#[serde(deny_unknown_fields)]
542struct HoldV1 {
543 expires_at_ms: u64,
544}
545
546#[derive(Serialize, Deserialize)]
547#[serde(deny_unknown_fields)]
548struct HolderV1 {
549 seq: u64,
550 #[serde(with = "hash_json")]
551 op_id: Hash,
552}
553
554#[derive(Serialize, Deserialize)]
555#[serde(deny_unknown_fields)]
556struct BlockV1 {
557 reason: String,
558 blocked_at_ms: u64,
559}
560
561#[derive(Serialize, Deserialize)]
562#[serde(deny_unknown_fields)]
563struct ObjectStateV1 {
564 seq: u64,
565 changed_at_ms: u64,
566 holders: u64,
567 deleting: bool,
568}
569
570const CODES: [Code; 16] = [
572 Code::Canceled,
573 Code::Unknown,
574 Code::InvalidArgument,
575 Code::DeadlineExceeded,
576 Code::NotFound,
577 Code::AlreadyExists,
578 Code::PermissionDenied,
579 Code::ResourceExhausted,
580 Code::FailedPrecondition,
581 Code::Aborted,
582 Code::OutOfRange,
583 Code::Unimplemented,
584 Code::Internal,
585 Code::Unavailable,
586 Code::DataLoss,
587 Code::Unauthenticated,
588];
589
590fn corrupt(what: &'static str) -> StoreError {
591 StoreError::Corrupt(what.into())
592}
593
594fn encode_json<T: Serialize>(value: &T) -> Value {
595 let mut out = vec![CODEC_V1];
596 serde_json::to_writer(&mut out, value).expect("codec DTOs always serialize");
597 Value::new(out)
598}
599
600fn decode_json<'a, T: Deserialize<'a>>(
601 value: &'a Value,
602 what: &'static str,
603) -> Result<T, StoreError> {
604 match value.as_bytes().split_first() {
605 Some((&CODEC_V1, body)) => serde_json::from_slice(body).map_err(|_| corrupt(what)),
606 _ => Err(corrupt("unknown codec version")),
607 }
608}
609
610fn hash_from(hex: &str) -> Result<Hash, StoreError> {
611 from_hex(hex).map_err(|_| corrupt("bad hash"))
612}
613
614#[must_use]
616pub fn encode_namespace_record(record: &NamespaceRecord) -> Value {
617 encode_json(record)
618}
619
620pub fn decode_namespace_record(value: &Value) -> Result<NamespaceRecord, StoreError> {
623 let record: NamespaceRecord = decode_json(value, "bad namespace record")?;
624 if record.config_version == 0 {
625 return Err(corrupt("namespace configuration version is zero"));
626 }
627 Ok(record)
628}
629
630#[must_use]
632pub fn encode_repo_record(record: &RepoRecord) -> Value {
633 encode_json(record)
634}
635
636pub fn decode_repo_record(value: &Value) -> Result<RepoRecord, StoreError> {
638 decode_json(value, "bad repo record")
639}
640
641#[must_use]
643pub fn encode_repo_visibility(row: &RepoVisibilityV1) -> Value {
644 encode_json(row)
645}
646
647pub fn decode_repo_visibility(value: &Value) -> Result<RepoVisibilityV1, StoreError> {
650 let row: RepoVisibilityV1 = decode_json(value, "bad repo visibility")?;
651 if let Some(id) = &row.last_statement_id
652 && !is_hex(id, 32)
653 {
654 return Err(corrupt("bad statement id"));
655 }
656 Ok(row)
657}
658
659#[must_use]
661pub fn encode_epoch_lease(lease: &EpochLease) -> Value {
662 encode_json(lease)
663}
664
665pub fn decode_epoch_lease(value: &Value) -> Result<EpochLease, StoreError> {
667 let lease: EpochLease = decode_json(value, "bad epoch lease")?;
668 if lease.authority_ready.is_some() && lease.authority_generation.is_none() {
669 return Err(corrupt("authority ready lease missing generation"));
670 }
671 if lease.config_version == 0 {
672 return Err(corrupt("lease configuration version is zero"));
673 }
674 Ok(lease)
675}
676
677#[must_use]
679pub fn encode_leased_shard(lease: &LeasedShard) -> Value {
680 encode_json(lease)
681}
682
683pub fn decode_leased_shard(value: &Value) -> Result<LeasedShard, StoreError> {
685 decode_json(value, "bad leased shard")
686}
687
688#[must_use]
690pub fn encode_lease_recovery(recovery: &LeaseRecovery) -> Value {
691 encode_json(recovery)
692}
693
694pub fn decode_lease_recovery(value: &Value) -> Result<LeaseRecovery, StoreError> {
696 let mode: LeaseRecovery = decode_json(value, "bad lease recovery")?;
697 if mode.authority_fence == Some(false)
698 || mode.activation_only == Some(false)
699 || ((mode.authority_ready.is_some() || mode.activation_only == Some(true))
700 && mode.authority_fence != Some(true))
701 {
702 return Err(corrupt("invalid authority activation marker"));
703 }
704 Ok(mode)
705}
706
707#[must_use]
709pub fn encode_backup_state(state: &BackupStateV1) -> Value {
710 encode_json(state)
711}
712
713pub fn decode_backup_state(value: &Value) -> Result<BackupStateV1, StoreError> {
715 decode_json(value, "bad backup state")
716}
717
718pub fn validate_ticket(ticket: &TicketV1) -> Result<(), StoreError> {
720 let ttl = ticket.expires_at_ms.checked_sub(ticket.created_at_ms);
721 if !is_served_ref_name(&ticket.ref_name)
722 || !validate_reservation_id(&ticket.reservation_id)
723 || ticket.bytes == 0
724 || ticket.part_size < MIN_PART_SIZE
725 || !ticket.part_size.is_power_of_two()
726 || !matches!(ttl, Some(1..604_800_000))
727 {
728 return Err(corrupt("invalid ticket"));
729 }
730 Ok(())
731}
732
733#[must_use]
735pub fn encode_ticket(ticket: &TicketV1) -> Value {
736 encode_json(ticket)
737}
738
739pub fn decode_ticket(value: &Value) -> Result<TicketV1, StoreError> {
741 check_value_limit(value)?;
742 let ticket = decode_json(value, "bad ticket")?;
743 validate_ticket(&ticket)?;
744 Ok(ticket)
745}
746
747#[must_use]
749pub fn encode_reservation(reservation: &ReservationV1) -> Value {
750 encode_json(reservation)
751}
752
753pub fn decode_reservation(value: &Value) -> Result<ReservationV1, StoreError> {
755 check_value_limit(value)?;
756 let reservation = decode_json(value, "bad reservation")?;
757 let repository = match &reservation {
758 ReservationV1::Ticketed { .. } => return Ok(reservation),
759 ReservationV1::Pending {
760 repository,
761 created_at_ms,
762 reconcile_at_ms,
763 ..
764 } => {
765 if reconcile_at_ms < created_at_ms {
766 return Err(corrupt("invalid pending deadline"));
767 }
768 repository
769 }
770 ReservationV1::Committed {
771 repository, refs, ..
772 } => {
773 for r in refs {
774 if !is_served_ref_name(&r.name) || r.deleted != r.new.is_none() {
775 return Err(corrupt("invalid outcome ref"));
776 }
777 }
778 repository
779 }
780 ReservationV1::Aborted {
781 repository, detail, ..
782 } => {
783 if detail.len() > 512 {
784 return Err(corrupt("outcome detail exceeds 512 bytes"));
785 }
786 repository
787 }
788 ReservationV1::Expired { repository, .. }
789 | ReservationV1::ReadServed { repository, .. } => repository,
790 };
791 RepositoryIdentity::parse_bare_allowed(repository)
792 .map_err(|_| corrupt("bad outcome repository"))?;
793 Ok(reservation)
794}
795
796fn check_value_limit(value: &Value) -> Result<(), StoreError> {
797 if value.as_bytes().len() > MAX_VALUE_BYTES {
798 return Err(corrupt("value exceeds MAX_VALUE_BYTES"));
799 }
800 Ok(())
801}
802
803fn hex_bytes(hex: &str) -> Result<Vec<u8>, StoreError> {
804 if !hex.len().is_multiple_of(2) {
805 return Err(corrupt("bad hex bytes"));
806 }
807 let digit = |b| match b {
808 b'0'..=b'9' => Some(b - b'0'),
809 b'a'..=b'f' => Some(b - b'a' + 10),
810 _ => None,
811 };
812 hex.as_bytes()
813 .chunks_exact(2)
814 .map(|pair| {
815 Ok(digit(pair[0]).ok_or_else(|| corrupt("bad hex bytes"))? * 16
816 + digit(pair[1]).ok_or_else(|| corrupt("bad hex bytes"))?)
817 })
818 .collect()
819}
820
821pub fn encode_relay(relay: &RelayV1) -> Result<Value, StoreError> {
823 let put_keys: std::collections::BTreeSet<_> = relay.puts.iter().map(|(key, _)| key).collect();
824 let mut delete_keys = std::collections::BTreeSet::new();
825 for key in &relay.deletes {
826 if key.as_bytes().len() > MAX_KEY_BYTES
827 || put_keys.contains(key)
828 || !delete_keys.insert(key)
829 {
830 return Err(StoreError::Invalid("invalid relay delete key".into()));
831 }
832 }
833 let value = encode_json(&RelayDtoV1 {
834 at_ms: relay.at_ms,
835 target: to_hex_bytes(&relay.target.encode()?),
836 puts: relay
837 .puts
838 .iter()
839 .map(|(key, value)| (to_hex_bytes(key.as_bytes()), to_hex_bytes(value.as_bytes())))
840 .collect(),
841 deletes: relay
842 .deletes
843 .iter()
844 .map(|key| to_hex_bytes(key.as_bytes()))
845 .collect(),
846 });
847 if value.as_bytes().len() > MAX_VALUE_BYTES {
848 return Err(StoreError::Invalid("relay exceeds MAX_VALUE_BYTES".into()));
849 }
850 for (key, value) in &relay.puts {
851 if key.as_bytes().len() > MAX_KEY_BYTES || value.as_bytes().len() > MAX_VALUE_BYTES {
852 return Err(StoreError::Invalid("invalid relay upsert size".into()));
853 }
854 }
855 Ok(value)
856}
857
858pub fn decode_relay(value: &Value) -> Result<RelayV1, StoreError> {
860 check_value_limit(value)?;
861 let dto: RelayDtoV1 = decode_json(value, "bad relay")?;
862 let target = Partition::decode(&hex_bytes(&dto.target)?)?;
863 let puts: Vec<(Key, Value)> = dto
864 .puts
865 .into_iter()
866 .map(|(key, value)| {
867 let key = hex_bytes(&key)?;
868 let value = hex_bytes(&value)?;
869 if key.len() > MAX_KEY_BYTES || value.len() > MAX_VALUE_BYTES {
870 return Err(corrupt("invalid relay upsert size"));
871 }
872 Ok((Key::new(key), Value::new(value)))
873 })
874 .collect::<Result<_, StoreError>>()?;
875 let put_keys: std::collections::BTreeSet<_> = puts.iter().map(|(key, _)| key).collect();
876 let mut delete_keys = std::collections::BTreeSet::new();
877 let mut deletes = Vec::with_capacity(dto.deletes.len());
878 for encoded in dto.deletes {
879 let key = Key::new(hex_bytes(&encoded)?);
880 if key.as_bytes().len() > MAX_KEY_BYTES
881 || put_keys.contains(&key)
882 || !delete_keys.insert(key.clone())
883 {
884 return Err(corrupt("invalid relay delete key"));
885 }
886 deletes.push(key);
887 }
888 Ok(RelayV1 {
889 at_ms: dto.at_ms,
890 target,
891 puts,
892 deletes,
893 })
894}
895
896pub fn encode_object_index(object: &Hash, row: &IndexValue) -> Result<Value, StoreError> {
900 row.validate(object)?;
901 let mut bytes = Vec::with_capacity(63);
902 bytes.push(CODEC_V1);
903 bytes.extend_from_slice(&row.frame_offset.to_be_bytes());
904 bytes.extend_from_slice(&row.frame_length.to_be_bytes());
905 bytes.push(row.wire_type);
906 bytes.extend_from_slice(&row.decoded_size.to_be_bytes());
907 bytes.extend_from_slice(&row.chain_depth.to_be_bytes());
908 match row.delta_base {
909 Some(base) => {
910 bytes.push(1);
911 bytes.extend_from_slice(&base);
912 }
913 None => bytes.push(0),
914 }
915 Ok(Value::new(bytes))
916}
917
918pub fn decode_object_index(object: &Hash, value: &Value) -> Result<IndexValue, StoreError> {
920 let bytes = value.as_bytes();
921 if !matches!(bytes.len(), 31 | 63) || bytes[0] != CODEC_V1 {
922 return Err(StoreError::Corrupt("bad object index value".into()));
923 }
924 let base = match (bytes[30], bytes.len()) {
925 (0, 31) => None,
926 (1, 63) => Some(
927 bytes[31..63]
928 .try_into()
929 .map_err(|_| StoreError::Corrupt("bad object index base".into()))?,
930 ),
931 _ => return Err(StoreError::Corrupt("bad object index base flag".into())),
932 };
933 let row = IndexValue {
934 frame_offset: u64::from_be_bytes(
935 bytes[1..9]
936 .try_into()
937 .map_err(|_| StoreError::Corrupt("bad object index offset".into()))?,
938 ),
939 frame_length: u64::from_be_bytes(
940 bytes[9..17]
941 .try_into()
942 .map_err(|_| StoreError::Corrupt("bad object index length".into()))?,
943 ),
944 wire_type: bytes[17],
945 decoded_size: u64::from_be_bytes(
946 bytes[18..26]
947 .try_into()
948 .map_err(|_| StoreError::Corrupt("bad object index size".into()))?,
949 ),
950 chain_depth: u32::from_be_bytes(
951 bytes[26..30]
952 .try_into()
953 .map_err(|_| StoreError::Corrupt("bad object index depth".into()))?,
954 ),
955 delta_base: base,
956 };
957 row.validate(object)
958 .map_err(|_| StoreError::Corrupt("invalid object index value".into()))?;
959 Ok(row)
960}
961
962fn relay_scan_invalid(scan: &RelayScanV1) -> Option<&'static str> {
963 if scan.cursor > scan.cycle_end {
964 Some("relay scan cursor exceeds cycle end")
965 } else if scan.blocked.len() > MAX_BLOCKED_TARGETS {
966 Some("relay scan exceeds MAX_BLOCKED_TARGETS")
967 } else if scan.blocked.windows(2).any(|pair| pair[0] >= pair[1]) {
968 Some("relay scan targets are not sorted and unique")
969 } else {
970 None
971 }
972}
973
974pub fn encode_relay_scan(scan: &RelayScanV1) -> Result<Value, StoreError> {
976 if let Some(message) = relay_scan_invalid(scan) {
977 return Err(StoreError::Invalid(message.into()));
978 }
979 let blocked = scan
980 .blocked
981 .iter()
982 .map(|target| Ok(to_hex_bytes(&target.encode()?)))
983 .collect::<Result<_, StoreError>>()?;
984 let value = encode_json(&RelayScanDtoV1 {
985 cycle_end: scan.cycle_end,
986 cursor: scan.cursor,
987 blocked,
988 });
989 if value.as_bytes().len() > MAX_VALUE_BYTES {
990 return Err(StoreError::Invalid(
991 "relay scan exceeds MAX_VALUE_BYTES".into(),
992 ));
993 }
994 Ok(value)
995}
996
997pub fn decode_relay_scan(value: &Value) -> Result<RelayScanV1, StoreError> {
999 check_value_limit(value)?;
1000 let dto: RelayScanDtoV1 = decode_json(value, "bad relay scan")?;
1001 if dto.blocked.len() > MAX_BLOCKED_TARGETS {
1002 return Err(corrupt("relay scan exceeds MAX_BLOCKED_TARGETS"));
1003 }
1004 let blocked = dto
1005 .blocked
1006 .into_iter()
1007 .map(|target| Partition::decode(&hex_bytes(&target)?))
1008 .collect::<Result<_, _>>()?;
1009 let scan = RelayScanV1 {
1010 cycle_end: dto.cycle_end,
1011 cursor: dto.cursor,
1012 blocked,
1013 };
1014 if let Some(message) = relay_scan_invalid(&scan) {
1015 return Err(corrupt(message));
1016 }
1017 Ok(scan)
1018}
1019
1020#[must_use]
1022pub fn encode_backlog(backlog: &Backlog) -> Value {
1023 encode_json(backlog)
1024}
1025
1026pub fn decode_backlog(value: &Value) -> Result<Backlog, StoreError> {
1028 check_value_limit(value)?;
1029 let backlog: Backlog = decode_json(value, "bad outcome backlog")?;
1030 if (backlog.rows == 0) != (backlog.bytes == 0) {
1031 return Err(corrupt("inconsistent outcome backlog"));
1032 }
1033 Ok(backlog)
1034}
1035
1036#[must_use]
1038pub fn encode_replay_record(record: &ReplayRecord) -> Value {
1039 let state = match &record.state {
1040 ReplayState::InFlight { resumable } => StateV1::InFlight {
1041 resumable: *resumable,
1042 },
1043 ReplayState::Committed(result) => StateV1::Committed {
1044 result: match result {
1045 StoredResult::UpdateRef(UpdateRefResult::Committed) => ResultV1::UpdateRefCommitted,
1046 StoredResult::UpdateRef(UpdateRefResult::Conflict { current }) => {
1047 ResultV1::UpdateRefConflict {
1048 current: current.as_ref().map(to_hex),
1049 }
1050 }
1051 StoredResult::AdvanceRefs(AdvanceOutcome::Committed) => ResultV1::AdvanceCommitted,
1052 StoredResult::AdvanceRefs(AdvanceOutcome::HeadConflict) => {
1053 ResultV1::AdvanceHeadConflict
1054 }
1055 StoredResult::AdvanceRefs(AdvanceOutcome::PackmapConflict) => {
1056 ResultV1::AdvancePackmapConflict
1057 }
1058 StoredResult::BeginUpload(BeginUploadResult::AlreadyPresent) => {
1059 ResultV1::BeginUploadAlreadyPresent
1060 }
1061 StoredResult::BeginUpload(BeginUploadResult::Ticket {
1062 id,
1063 part_size,
1064 expires_at_ms,
1065 token,
1066 }) => ResultV1::BeginUploadTicket {
1067 id: to_hex(id),
1068 part_size: *part_size,
1069 expires_at_ms: *expires_at_ms,
1070 token_hex: to_hex_bytes(token),
1071 },
1072 StoredResult::UploadPack => ResultV1::UploadPack,
1073 StoredResult::RepoVisibility => ResultV1::RepoVisibility,
1074 StoredResult::Rejected(r) => ResultV1::Rejected {
1075 code: r.code().as_str().to_owned(),
1076 message: r.message().to_owned(),
1077 },
1078 },
1079 },
1080 };
1081 encode_json(&RecordV1 {
1082 fingerprint: to_hex(&record.fingerprint),
1083 expires_at_ms: record.expires_at_ms,
1084 state,
1085 })
1086}
1087
1088pub fn decode_replay_record(value: &Value) -> Result<ReplayRecord, StoreError> {
1090 let dto: RecordV1 = decode_json(value, "bad replay record")?;
1091 let state = match dto.state {
1092 StateV1::InFlight { resumable } => ReplayState::InFlight { resumable },
1093 StateV1::Committed { result } => ReplayState::Committed(match result {
1094 ResultV1::UpdateRefCommitted => StoredResult::UpdateRef(UpdateRefResult::Committed),
1095 ResultV1::UpdateRefConflict { current } => {
1096 StoredResult::UpdateRef(UpdateRefResult::Conflict {
1097 current: current.as_deref().map(hash_from).transpose()?,
1098 })
1099 }
1100 ResultV1::AdvanceCommitted => StoredResult::AdvanceRefs(AdvanceOutcome::Committed),
1101 ResultV1::AdvanceHeadConflict => {
1102 StoredResult::AdvanceRefs(AdvanceOutcome::HeadConflict)
1103 }
1104 ResultV1::AdvancePackmapConflict => {
1105 StoredResult::AdvanceRefs(AdvanceOutcome::PackmapConflict)
1106 }
1107 ResultV1::BeginUploadAlreadyPresent => {
1108 StoredResult::BeginUpload(BeginUploadResult::AlreadyPresent)
1109 }
1110 ResultV1::BeginUploadTicket {
1111 id,
1112 part_size,
1113 expires_at_ms,
1114 token_hex,
1115 } => {
1116 if part_size < mkit_core::upload_parts::MIN_PART_SIZE
1117 || !part_size.is_power_of_two()
1118 || token_hex.is_empty()
1119 {
1120 return Err(corrupt("invalid stored ticket result"));
1121 }
1122 StoredResult::BeginUpload(BeginUploadResult::Ticket {
1123 id: hash_from(&id)?,
1124 part_size,
1125 expires_at_ms,
1126 token: hex_bytes(&token_hex)?,
1127 })
1128 }
1129 ResultV1::UploadPack => StoredResult::UploadPack,
1130 ResultV1::RepoVisibility => StoredResult::RepoVisibility,
1131 ResultV1::Rejected { code, message } => {
1132 let code = CODES
1133 .into_iter()
1134 .find(|c| c.as_str() == code)
1135 .ok_or_else(|| corrupt("unknown code"))?;
1136 StoredResult::Rejected(
1137 StoredRejection::new(code, message)
1138 .ok_or_else(|| corrupt("stored rejection code is not final"))?,
1139 )
1140 }
1141 }),
1142 };
1143 Ok(ReplayRecord {
1144 fingerprint: hash_from(&dto.fingerprint)?,
1145 expires_at_ms: dto.expires_at_ms,
1146 state,
1147 })
1148}
1149
1150#[must_use]
1152pub fn encode_quota_state(state: &QuotaState) -> Value {
1153 encode_json(&QuotaV1 {
1154 window_start: state.window_start,
1155 ops: state.ops,
1156 bytes: state.bytes,
1157 })
1158}
1159
1160pub fn decode_quota_state(value: &Value) -> Result<QuotaState, StoreError> {
1162 let dto: QuotaV1 = decode_json(value, "bad quota state")?;
1163 Ok(QuotaState {
1164 window_start: dto.window_start,
1165 ops: dto.ops,
1166 bytes: dto.bytes,
1167 })
1168}
1169
1170#[must_use]
1172pub fn encode_namespace_usage(usage: NamespaceUsage) -> Value {
1173 Value::new([usage.ops.to_be_bytes(), usage.bytes.to_be_bytes()].concat())
1174}
1175
1176pub fn decode_namespace_usage(value: &Value) -> Result<NamespaceUsage, StoreError> {
1178 let (ops, bytes) = value
1179 .as_bytes()
1180 .split_first_chunk::<8>()
1181 .ok_or_else(|| corrupt("bad namespace usage"))?;
1182 let bytes: [u8; 8] = bytes
1183 .try_into()
1184 .map_err(|_| corrupt("bad namespace usage"))?;
1185 Ok(NamespaceUsage {
1186 ops: u64::from_be_bytes(*ops),
1187 bytes: u64::from_be_bytes(bytes),
1188 })
1189}
1190
1191#[must_use]
1193pub fn encode_namespace_view(view: NamespaceView) -> Value {
1194 Value::new(
1195 [
1196 view.total.ops.to_be_bytes(),
1197 view.total.bytes.to_be_bytes(),
1198 view.pushed.ops.to_be_bytes(),
1199 view.pushed.bytes.to_be_bytes(),
1200 view.observed_at_ms.to_be_bytes(),
1201 ]
1202 .concat(),
1203 )
1204}
1205
1206pub fn decode_namespace_view(value: &Value) -> Result<NamespaceView, StoreError> {
1208 let bytes: [u8; 40] = value
1209 .as_bytes()
1210 .try_into()
1211 .map_err(|_| corrupt("bad namespace view"))?;
1212 let word = |i| -> Result<u64, StoreError> {
1213 let chunk: [u8; 8] = bytes
1214 .get(i..i + 8)
1215 .ok_or_else(|| corrupt("bad namespace view"))?
1216 .try_into()
1217 .map_err(|_| corrupt("bad namespace view"))?;
1218 Ok(u64::from_be_bytes(chunk))
1219 };
1220 let view = NamespaceView {
1221 total: NamespaceUsage {
1222 ops: word(0)?,
1223 bytes: word(8)?,
1224 },
1225 pushed: NamespaceUsage {
1226 ops: word(16)?,
1227 bytes: word(24)?,
1228 },
1229 observed_at_ms: word(32)?,
1230 };
1231 if view.total.delta_from(view.pushed).is_none() {
1232 return Err(corrupt("namespace view exceeds total"));
1233 }
1234 Ok(view)
1235}
1236
1237#[must_use]
1239pub fn encode_hold(expires_at_ms: u64) -> Value {
1240 encode_json(&HoldV1 { expires_at_ms })
1241}
1242
1243pub fn decode_hold(value: &Value) -> Result<u64, StoreError> {
1245 let dto: HoldV1 = decode_json(value, "bad hold")?;
1246 Ok(dto.expires_at_ms)
1247}
1248
1249#[must_use]
1251pub fn encode_holder(record: &HolderRecord) -> Value {
1252 encode_json(&HolderV1 {
1253 seq: record.seq,
1254 op_id: record.op_id,
1255 })
1256}
1257
1258pub fn decode_holder(value: &Value) -> Result<HolderRecord, StoreError> {
1260 let dto: HolderV1 = decode_json(value, "bad holder")?;
1261 Ok(HolderRecord::new(dto.seq, dto.op_id))
1262}
1263
1264#[must_use]
1266pub fn encode_block_entry(entry: &BlockEntry) -> Value {
1267 encode_json(&BlockV1 {
1268 reason: entry.reason.clone(),
1269 blocked_at_ms: entry.blocked_at_ms,
1270 })
1271}
1272
1273pub fn decode_block_entry(value: &Value) -> Result<BlockEntry, StoreError> {
1275 let dto: BlockV1 = decode_json(value, "bad blocklist entry")?;
1276 Ok(BlockEntry {
1277 reason: dto.reason,
1278 blocked_at_ms: dto.blocked_at_ms,
1279 })
1280}
1281
1282#[must_use]
1284pub fn encode_object_state(state: &ObjectState) -> Value {
1285 encode_json(&ObjectStateV1 {
1286 seq: state.seq,
1287 changed_at_ms: state.changed_at_ms,
1288 holders: state.holders,
1289 deleting: state.deleting,
1290 })
1291}
1292
1293pub fn decode_object_state(value: &Value) -> Result<ObjectState, StoreError> {
1295 let dto: ObjectStateV1 = decode_json(value, "bad object state")?;
1296 Ok(ObjectState {
1297 seq: dto.seq,
1298 changed_at_ms: dto.changed_at_ms,
1299 holders: dto.holders,
1300 deleting: dto.deleting,
1301 })
1302}
1303
1304#[must_use]
1306pub fn encode_ref_id(id: &Hash) -> Value {
1307 Value::new(id.to_vec())
1308}
1309
1310pub fn decode_ref_id(value: &Value) -> Result<Hash, StoreError> {
1312 Hash::try_from(value.as_bytes()).map_err(|_| corrupt("ref value is not 32 bytes"))
1313}
1314
1315#[must_use]
1317pub fn encode_u64(n: u64) -> Value {
1318 Value::new(n.to_be_bytes().to_vec())
1319}
1320
1321pub fn decode_u64(value: &Value) -> Result<u64, StoreError> {
1323 <[u8; 8]>::try_from(value.as_bytes())
1324 .map(u64::from_be_bytes)
1325 .map_err(|_| corrupt("integer is not 8 bytes"))
1326}
1327
1328#[must_use]
1330pub fn encode_u32(n: u32) -> Value {
1331 Value::new(n.to_be_bytes().to_vec())
1332}
1333
1334pub fn decode_u32(value: &Value) -> Result<u32, StoreError> {
1336 <[u8; 4]>::try_from(value.as_bytes())
1337 .map(u32::from_be_bytes)
1338 .map_err(|_| corrupt("integer is not 4 bytes"))
1339}
1340
1341#[cfg(test)]
1342mod tests {
1343 use super::*;
1344
1345 fn ticket_fixture() -> TicketV1 {
1346 TicketV1 {
1347 authority_generation: None,
1348 repo: RepoName::new("a").unwrap(),
1349 ref_name: "refs/heads/main".into(),
1350 signer: [0x11; 32],
1351 pack_id: [0x22; 32],
1352 bytes: 9,
1353 part_size: MIN_PART_SIZE,
1354 expires_at_ms: 24,
1355 created_at_ms: 1,
1356 reservation_id: "R-1:ok".into(),
1357 upload_session: None,
1358 }
1359 }
1360
1361 fn json_value(json: &serde_json::Value) -> Value {
1362 let mut bytes = vec![CODEC_V1];
1363 bytes.extend(serde_json::to_vec(json).unwrap());
1364 Value::new(bytes)
1365 }
1366
1367 #[test]
1368 fn ticket_codec_golden_roundtrip_and_rejections() {
1369 let ticket = ticket_fixture();
1370 let golden = format!(
1371 r#"{{"repo":"a","ref_name":"refs/heads/main","signer":"{}","pack_id":"{}","bytes":9,"part_size":8388608,"expires_at_ms":24,"created_at_ms":1,"reservation_id":"R-1:ok","upload_session":null}}"#,
1372 "11".repeat(32),
1373 "22".repeat(32)
1374 );
1375 assert_eq!(
1376 encode_ticket(&ticket).as_bytes(),
1377 [&[CODEC_V1][..], golden.as_bytes()].concat()
1378 );
1379 assert_eq!(decode_ticket(&encode_ticket(&ticket)).unwrap(), ticket);
1380 let mut session = ticket.clone();
1381 session.upload_session = Some(b"backend-session".to_vec());
1382 let encoded = encode_ticket(&session);
1383 let expected = golden.replace(
1384 "\"upload_session\":null",
1385 "\"upload_session\":\"6261636b656e642d73657373696f6e\"",
1386 );
1387 assert_eq!(
1388 encoded.as_bytes(),
1389 [&[CODEC_V1][..], expected.as_bytes()].concat()
1390 );
1391 assert_eq!(decode_ticket(&encode_ticket(&session)).unwrap(), session);
1392 let base = serde_json::to_value(&ticket).unwrap();
1393 for (field, bad) in [
1394 ("signer", serde_json::json!("bad hex")),
1395 ("pack_id", serde_json::json!("00")),
1396 ("bytes", serde_json::json!(0)),
1397 ("part_size", serde_json::json!(8_388_609)),
1398 ("part_size", serde_json::json!(1)),
1399 ("unknown", serde_json::json!(1)),
1400 ("upload_session", serde_json::json!("+0")),
1401 ("upload_session", serde_json::json!("0+")),
1402 ("upload_session", serde_json::json!("é0")),
1403 ("reservation_id", serde_json::json!("bad/id")),
1404 ("reservation_id", serde_json::json!("")),
1405 ("reservation_id", serde_json::json!("a".repeat(129))),
1406 ("repo", serde_json::json!("bad name")),
1407 ("ref_name", serde_json::json!("refs/heads/../b")),
1408 ("expires_at_ms", serde_json::json!(1)),
1409 ("expires_at_ms", serde_json::json!(604_800_001)),
1410 ("created_at_ms", serde_json::json!(25)),
1411 ("bytes", serde_json::json!(-1)),
1412 ] {
1413 let mut bad_json = base.clone();
1414 bad_json[field] = bad;
1415 assert!(
1416 matches!(
1417 decode_ticket(&json_value(&bad_json)),
1418 Err(StoreError::Corrupt(_))
1419 ),
1420 "{field}"
1421 );
1422 }
1423 let mut boundary = ticket;
1424 boundary.expires_at_ms = boundary.created_at_ms + 604_799_999;
1425 assert!(decode_ticket(&encode_ticket(&boundary)).is_ok());
1426 for bad in [
1427 Value::new(vec![]),
1428 Value::new(b"\x02{}".to_vec()),
1429 Value::new(b"\x01{".to_vec()),
1430 Value::new(vec![0; MAX_VALUE_BYTES + 1]),
1431 ] {
1432 assert!(decode_ticket(&bad).is_err());
1433 }
1434 }
1435
1436 #[test]
1437 fn reservation_codec_all_variants_golden_and_roundtrip() {
1438 let cases = vec![
1439 (ReservationV1::Pending { repository: "a".into(), created_at_ms: 1, reconcile_at_ms: 300_001, op: PendingOp::Write }, r#"{"state":"pending","repository":"a","created_at_ms":1,"reconcile_at_ms":300001,"op":"write"}"#.into()),
1440 (ReservationV1::Pending { repository: "a".into(), created_at_ms: 1, reconcile_at_ms: 60_001, op: PendingOp::Read }, r#"{"state":"pending","repository":"a","created_at_ms":1,"reconcile_at_ms":60001,"op":"read"}"#.into()),
1441 (ReservationV1::Ticketed { ticket_id: [0x11; 32] }, format!(r#"{{"state":"ticketed","ticket_id":"{}"}}"#, "11".repeat(32))),
1442 (ReservationV1::Committed { repository: "a".into(), occurred_at_ms: 7, bytes_stored: 9, new_to_repo: 8, new_to_store: 6, refs: vec![OutcomeRef { name: "refs/heads/main".into(), new: Some([0x22; 32]), deleted: false }, OutcomeRef { name: "refs/tags/v1".into(), new: None, deleted: true }] }, format!(r#"{{"state":"committed","repository":"a","occurred_at_ms":7,"bytes_stored":9,"new_to_repo":8,"new_to_store":6,"refs":[{{"name":"refs/heads/main","new":"{}","deleted":false}},{{"name":"refs/tags/v1","new":null,"deleted":true}}]}}"#, "22".repeat(32))),
1443 (ReservationV1::Aborted { repository: "a".into(), occurred_at_ms: 7, reason: AbortReason::Abandoned, detail: "gone".into() }, r#"{"state":"aborted","repository":"a","occurred_at_ms":7,"reason":"ABANDONED","detail":"gone"}"#.into()),
1444 (ReservationV1::Expired { repository: "a".into(), occurred_at_ms: 7 }, r#"{"state":"expired","repository":"a","occurred_at_ms":7}"#.into()),
1445 (ReservationV1::ReadServed { repository: "a".into(), occurred_at_ms: 7, object: [0x33; 32], bytes_served: 9 }, format!(r#"{{"state":"read_served","repository":"a","occurred_at_ms":7,"object":"{}","bytes_served":9}}"#, "33".repeat(32))),
1446 ];
1447 for (row, golden) in cases {
1448 let value = encode_reservation(&row);
1449 assert_eq!(
1450 value.as_bytes(),
1451 [&[CODEC_V1][..], golden.as_bytes()].concat()
1452 );
1453 assert_eq!(decode_reservation(&value).unwrap(), row);
1454 let mut json = serde_json::to_value(&row).unwrap();
1455 json["extra"] = serde_json::json!(1);
1456 assert!(decode_reservation(&json_value(&json)).is_err());
1457 }
1458 for (reason, name) in [
1459 (AbortReason::Unspecified, "UNSPECIFIED"),
1460 (AbortReason::RefConflict, "REF_CONFLICT"),
1461 (AbortReason::EpochMismatch, "EPOCH_MISMATCH"),
1462 (AbortReason::PackMissing, "PACK_MISSING"),
1463 (AbortReason::ReplayRace, "REPLAY_RACE"),
1464 (AbortReason::Internal, "INTERNAL"),
1465 (AbortReason::Abandoned, "ABANDONED"),
1466 ] {
1467 let row = ReservationV1::Aborted {
1468 repository: "a".into(),
1469 occurred_at_ms: 7,
1470 reason,
1471 detail: String::new(),
1472 };
1473 let value = encode_reservation(&row);
1474 assert_eq!(value.as_bytes(), format!("\x01{{\"state\":\"aborted\",\"repository\":\"a\",\"occurred_at_ms\":7,\"reason\":\"{name}\",\"detail\":\"\"}}").as_bytes());
1475 assert_eq!(decode_reservation(&value).unwrap(), row);
1476 }
1477 let full_repo = format!(
1478 "ed25519-{}/{}",
1479 "ab".repeat(32),
1480 "r".repeat(mkit_core::repo_identity::MAX_NAME_LEN)
1481 );
1482 let row = ReservationV1::Expired {
1483 repository: full_repo,
1484 occurred_at_ms: 7,
1485 };
1486 assert_eq!(decode_reservation(&encode_reservation(&row)).unwrap(), row);
1487 }
1488
1489 #[test]
1490 fn reservation_codec_rejects_malformed_states_and_outcomes() {
1491 for json in [
1492 serde_json::json!({"state":"unknown"}),
1493 serde_json::json!({"state":"pending"}),
1494 serde_json::json!({"state":"ticketed","ticket_id":"nope"}),
1495 serde_json::json!({"state":"expired","repository":"bad name","occurred_at_ms":7}),
1496 serde_json::json!({"state":"aborted","repository":"a","occurred_at_ms":7,"reason":"UNKNOWN","detail":""}),
1497 serde_json::json!({"state":"aborted","repository":"a","occurred_at_ms":7,"reason":"INTERNAL","detail":"x".repeat(513)}),
1498 serde_json::json!({"state":"aborted","repository":"a","occurred_at_ms":7,"reason":"INTERNAL","detail":"é".repeat(257)}),
1499 serde_json::json!({"state":"committed","repository":"a","occurred_at_ms":7,"bytes_stored":1,"new_to_repo":1,"new_to_store":0,"refs":[{"name":"bad","new":null,"deleted":true}]}),
1500 serde_json::json!({"state":"committed","repository":"a","occurred_at_ms":7,"bytes_stored":1,"new_to_repo":1,"new_to_store":0,"refs":[{"name":"refs/heads/a","new":null,"deleted":false}]}),
1501 serde_json::json!({"state":"committed","repository":"a","occurred_at_ms":7,"bytes_stored":1,"new_to_repo":1,"new_to_store":0,"refs":[{"name":"refs/heads/a","new":"bad","deleted":false}]}),
1502 serde_json::json!({"state":"committed","repository":"a","occurred_at_ms":7,"bytes_stored":1,"new_to_repo":1,"new_to_store":0,"refs":[{"name":"refs/heads/a","new":null,"deleted":true,"extra":1}]}),
1503 ] {
1504 assert!(matches!(
1505 decode_reservation(&json_value(&json)),
1506 Err(StoreError::Corrupt(_))
1507 ));
1508 }
1509 let boundary = ReservationV1::Aborted {
1510 repository: "a".into(),
1511 occurred_at_ms: 7,
1512 reason: AbortReason::Internal,
1513 detail: "é".repeat(256),
1514 };
1515 assert_eq!(
1516 decode_reservation(&encode_reservation(&boundary)).unwrap(),
1517 boundary
1518 );
1519 for value in [
1520 Value::new(vec![]),
1521 Value::new(b"\x02{}".to_vec()),
1522 Value::new(b"\x01{".to_vec()),
1523 ] {
1524 assert!(decode_reservation(&value).is_err());
1525 }
1526 }
1527
1528 #[test]
1529 fn relay_and_backlog_codec_golden_roundtrip_and_rejections() {
1530 let relay = RelayV1 {
1531 at_ms: 123,
1532 target: Partition::Namespace(crate::repo::NamespaceKey::deployment_default()),
1533 puts: vec![(Key::new(b"m\0a\0".to_vec()), Value::new(vec![]))],
1534 deletes: Vec::new(),
1535 };
1536 let encoded = encode_relay(&relay).unwrap();
1537 assert_eq!(
1538 encoded.as_bytes(),
1539 b"\x01{\"at_ms\":123,\"target\":\"6e726f6f7400\",\"puts\":[[\"6d006100\",\"\"]]}"
1540 );
1541 assert_eq!(decode_relay(&encoded).unwrap(), relay);
1542 let backlog = Backlog { rows: 3, bytes: 72 };
1543 let encoded = encode_backlog(&backlog);
1544 assert_eq!(encoded.as_bytes(), b"\x01{\"rows\":3,\"bytes\":72}");
1545 assert_eq!(decode_backlog(&encoded).unwrap(), backlog);
1546 assert_eq!(
1547 decode_backlog(&encode_backlog(&Backlog::default())).unwrap(),
1548 Backlog::default()
1549 );
1550 for json in [
1551 serde_json::json!({"target":"6e726f6f7400","puts":[]}),
1552 serde_json::json!({"at_ms":-1,"target":"6e726f6f7400","puts":[]}),
1553 serde_json::json!({"at_ms":123,"target":"bad","puts":[]}),
1554 serde_json::json!({"at_ms":123,"target":"zz","puts":[]}),
1555 serde_json::json!({"at_ms":123,"target":"00","puts":[]}),
1556 serde_json::json!({"at_ms":123,"target":"6e726f6f7400","puts":[["gg",""]]}),
1557 serde_json::json!({"at_ms":123,"target":"6e726f6f7400","puts":[],"extra":1}),
1558 serde_json::json!({"at_ms":123,"target":"6e726f6f7400","puts":[["00".repeat(MAX_KEY_BYTES + 1),""]]}),
1559 ] {
1560 assert!(decode_relay(&json_value(&json)).is_err());
1561 }
1562 for json in [
1563 serde_json::json!({"rows":-1,"bytes":1}),
1564 serde_json::json!({"rows":1,"bytes":0}),
1565 serde_json::json!({"rows":0,"bytes":1}),
1566 serde_json::json!({"rows":1,"bytes":1,"extra":1}),
1567 ] {
1568 assert!(decode_backlog(&json_value(&json)).is_err());
1569 }
1570 for value in [
1571 Value::new(vec![]),
1572 Value::new(b"\x02{}".to_vec()),
1573 Value::new(b"\x01{".to_vec()),
1574 Value::new(vec![0; MAX_VALUE_BYTES + 1]),
1575 ] {
1576 assert!(decode_backlog(&value).is_err());
1577 assert!(decode_relay(&value).is_err());
1578 }
1579 let oversized = RelayV1 {
1580 at_ms: 123,
1581 target: relay.target.clone(),
1582 puts: vec![(Key::new(vec![0; MAX_KEY_BYTES + 1]), Value::new(vec![]))],
1583 deletes: Vec::new(),
1584 };
1585 assert!(encode_relay(&oversized).is_err());
1586 let key = Key::new(b"x\0a\0refs/heads/main".as_slice());
1587 let mut deleted = relay.clone();
1588 deleted.deletes.push(key.clone());
1589 assert_eq!(
1590 decode_relay(&encode_relay(&deleted).unwrap()).unwrap(),
1591 deleted
1592 );
1593 deleted.puts.push((key.clone(), Value::default()));
1594 assert!(encode_relay(&deleted).is_err());
1595 deleted.puts.pop();
1596 deleted.deletes.push(key.clone());
1597 assert!(encode_relay(&deleted).is_err());
1598 deleted.deletes = vec![Key::new(vec![b'x'; MAX_KEY_BYTES + 1])];
1599 assert!(encode_relay(&deleted).is_err());
1600 for json in [
1601 serde_json::json!({"at_ms":123,"target":"6e726f6f7400","puts":[["780061",""]],"deletes":["780061"]}),
1602 serde_json::json!({"at_ms":123,"target":"6e726f6f7400","puts":[],"deletes":["78","78"]}),
1603 serde_json::json!({"at_ms":123,"target":"6e726f6f7400","puts":[],"deletes":["78"],"extra":1}),
1604 serde_json::json!({"at_ms":123,"target":"6e726f6f7400","puts":[],"deletes":["78".repeat(MAX_KEY_BYTES + 1)]}),
1605 ] {
1606 assert!(decode_relay(&json_value(&json)).is_err());
1607 }
1608 }
1609
1610 #[test]
1611 fn relay_scan_codec_golden_and_roundtrip() {
1612 let scan = RelayScanV1 {
1613 cycle_end: 17,
1614 cursor: 3,
1615 blocked: vec![
1616 Partition::Namespace(crate::repo::NamespaceKey::deployment_default()),
1617 Partition::ContentShard(7),
1618 ],
1619 };
1620 let encoded = encode_relay_scan(&scan).unwrap();
1621 assert_eq!(
1622 encoded.as_bytes(),
1623 b"\x01{\"cycle_end\":17,\"cursor\":3,\"blocked\":[\"6e726f6f7400\",\"733700\"]}"
1624 );
1625 assert_eq!(decode_relay_scan(&encoded).unwrap(), scan);
1626 for (cycle_end, cursor) in [(0, 0), (17, 0), (17, 17), (u64::MAX, u64::MAX)] {
1627 let scan = RelayScanV1 {
1628 cycle_end,
1629 cursor,
1630 blocked: vec![],
1631 };
1632 assert_eq!(
1633 decode_relay_scan(&encode_relay_scan(&scan).unwrap()).unwrap(),
1634 scan
1635 );
1636 }
1637 }
1638
1639 #[test]
1640 fn relay_scan_codec_rejects_invalid_progress_and_blocked_targets() {
1641 let scans = [
1642 RelayScanV1 {
1643 cycle_end: 1,
1644 cursor: 2,
1645 blocked: vec![],
1646 },
1647 RelayScanV1 {
1648 cycle_end: 1,
1649 cursor: 0,
1650 blocked: vec![Partition::ContentShard(2), Partition::ContentShard(1)],
1651 },
1652 RelayScanV1 {
1653 cycle_end: 1,
1654 cursor: 0,
1655 blocked: vec![Partition::ContentShard(1), Partition::ContentShard(1)],
1656 },
1657 RelayScanV1 {
1658 cycle_end: 1,
1659 cursor: 0,
1660 blocked: (0..=32u16).map(Partition::ContentShard).collect(),
1661 },
1662 RelayScanV1 {
1663 cycle_end: 1,
1664 cursor: 0,
1665 blocked: vec![Partition::Namespace(
1666 crate::repo::NamespaceKey::from_stored("bad\0ns".into()),
1667 )],
1668 },
1669 ];
1670 for scan in scans {
1671 assert!(matches!(
1672 encode_relay_scan(&scan),
1673 Err(StoreError::Invalid(_))
1674 ));
1675 }
1676 for json in [
1677 serde_json::json!({"cycle_end":1,"cursor":2,"blocked":[]}),
1678 serde_json::json!({"cycle_end":1,"cursor":0,"blocked":["733200","733100"]}),
1679 serde_json::json!({"cycle_end":1,"cursor":0,"blocked":["733100","733100"]}),
1680 serde_json::json!({"cycle_end":1,"cursor":0,"blocked":(0..=32u16).map(|n| to_hex_bytes(&Partition::ContentShard(n).encode().unwrap())).collect::<Vec<_>>()}),
1681 serde_json::json!({"cycle_end":1,"cursor":0,"blocked":["gg"]}),
1682 serde_json::json!({"cycle_end":1,"cursor":0,"blocked":["0"]}),
1683 serde_json::json!({"cycle_end":1,"cursor":0,"blocked":["00"]}),
1684 serde_json::json!({"cycle_end":-1,"cursor":0,"blocked":[]}),
1685 serde_json::json!({"cycle_end":1,"cursor":0}),
1686 serde_json::json!({"cycle_end":1,"cursor":0,"blocked":[],"extra":1}),
1687 ] {
1688 assert!(matches!(
1689 decode_relay_scan(&json_value(&json)),
1690 Err(StoreError::Corrupt(_))
1691 ));
1692 }
1693 for value in [
1694 Value::new(vec![]),
1695 Value::new(b"\x02{}".to_vec()),
1696 Value::new(b"\x01{".to_vec()),
1697 Value::new(vec![0; MAX_VALUE_BYTES + 1]),
1698 ] {
1699 assert!(matches!(
1700 decode_relay_scan(&value),
1701 Err(StoreError::Corrupt(_))
1702 ));
1703 }
1704 }
1705
1706 #[test]
1707 fn relay_scan_32_maximum_partitions_fit_the_value_limit() {
1708 use crate::refs::MAX_REF_NAME_BYTES;
1709 use crate::repo::{MAX_REPO_NAME_BYTES, NamespaceKey};
1710
1711 assert_eq!(MAX_BLOCKED_TARGETS, 32);
1712 let namespace = NamespaceKey::from_stored(format!("ed25519-{}", "a".repeat(64)));
1713 let repo = RepoName::new("r".repeat(MAX_REPO_NAME_BYTES)).unwrap();
1714 let base = format!(
1715 "refs/heads/{}",
1716 "a".repeat(MAX_REF_NAME_BYTES - "refs/heads/".len() - 2)
1717 );
1718 let blocked = (0..32)
1719 .map(|n| {
1720 let shard_ref = format!("{base}{n:02}");
1721 assert_eq!(shard_ref.len(), MAX_REF_NAME_BYTES);
1722 assert!(crate::refs::validate_ref_name(&shard_ref));
1723 Partition::Ref {
1724 ns: namespace.clone(),
1725 repo: repo.clone(),
1726 shard_ref,
1727 }
1728 })
1729 .collect();
1730 let scan = RelayScanV1 {
1731 cycle_end: u64::MAX,
1732 cursor: u64::MAX,
1733 blocked,
1734 };
1735 let encoded = encode_relay_scan(&scan).unwrap();
1736 assert!(encoded.as_bytes().len() < MAX_VALUE_BYTES);
1737 assert_eq!(decode_relay_scan(&encoded).unwrap(), scan);
1738
1739 let oversized = RelayScanV1 {
1740 cycle_end: 1,
1741 cursor: 0,
1742 blocked: vec![Partition::Namespace(NamespaceKey::from_stored(
1743 "a".repeat(MAX_VALUE_BYTES),
1744 ))],
1745 };
1746 assert!(matches!(
1747 encode_relay_scan(&oversized),
1748 Err(StoreError::Invalid(_))
1749 ));
1750 }
1751
1752 fn records() -> Vec<ReplayRecord> {
1753 let results = [
1754 StoredResult::UpdateRef(UpdateRefResult::Committed),
1755 StoredResult::UpdateRef(UpdateRefResult::Conflict { current: None }),
1756 StoredResult::UpdateRef(UpdateRefResult::Conflict {
1757 current: Some([3; 32]),
1758 }),
1759 StoredResult::AdvanceRefs(AdvanceOutcome::Committed),
1760 StoredResult::AdvanceRefs(AdvanceOutcome::HeadConflict),
1761 StoredResult::AdvanceRefs(AdvanceOutcome::PackmapConflict),
1762 StoredResult::UploadPack,
1763 StoredResult::RepoVisibility,
1764 StoredResult::Rejected(StoredRejection::new(Code::PermissionDenied, "no").unwrap()),
1765 ];
1766 let mut states: Vec<_> = results.into_iter().map(ReplayState::Committed).collect();
1767 states.push(ReplayState::InFlight { resumable: true });
1768 states.push(ReplayState::InFlight { resumable: false });
1769 states
1770 .into_iter()
1771 .map(|state| ReplayRecord {
1772 fingerprint: [9; 32],
1773 expires_at_ms: -5,
1774 state,
1775 })
1776 .collect()
1777 }
1778
1779 #[test]
1780 fn codec_roundtrip_every_value_type() {
1781 for record in records() {
1782 let value = encode_replay_record(&record);
1783 assert_eq!(value.as_bytes()[0], CODEC_V1);
1784 assert_eq!(decode_replay_record(&value).unwrap(), record);
1785 }
1786 let quota = QuotaState {
1787 window_start: 1_700_000_000_000,
1788 ops: 3,
1789 bytes: u64::MAX,
1790 };
1791 assert_eq!(
1792 decode_quota_state(&encode_quota_state("a)).unwrap(),
1793 quota
1794 );
1795 assert_eq!(decode_ref_id(&encode_ref_id(&[4; 32])).unwrap(), [4; 32]);
1796 assert_eq!(decode_u64(&encode_u64(u64::MAX - 1)).unwrap(), u64::MAX - 1);
1797 assert_eq!(
1798 decode_u32(&encode_u32(1)).unwrap().to_be_bytes(),
1799 [0, 0, 0, 1]
1800 );
1801 for code in CODES {
1802 assert_eq!(
1803 CODES.iter().filter(|c| c.as_str() == code.as_str()).count(),
1804 1
1805 );
1806 }
1807 }
1808
1809 #[test]
1810 fn codec_golden_bytes() {
1811 let head = format!(
1812 "\x01{{\"fingerprint\":\"{}\",\"expires_at_ms\":-5,\"state\":",
1813 "09".repeat(32)
1814 );
1815 let committed = |result: &str| format!("{{\"state\":\"committed\",\"result\":{result}}}");
1816 let states = [
1817 committed(r#"{"kind":"update_ref_committed"}"#),
1818 committed(r#"{"kind":"update_ref_conflict","current":null}"#),
1819 committed(&format!(
1820 r#"{{"kind":"update_ref_conflict","current":"{}"}}"#,
1821 "03".repeat(32)
1822 )),
1823 committed(r#"{"kind":"advance_committed"}"#),
1824 committed(r#"{"kind":"advance_head_conflict"}"#),
1825 committed(r#"{"kind":"advance_packmap_conflict"}"#),
1826 committed(r#"{"kind":"upload_pack"}"#),
1827 committed(r#"{"kind":"repo_visibility"}"#),
1828 committed(r#"{"kind":"rejected","code":"permission_denied","message":"no"}"#),
1829 r#"{"state":"in_flight","resumable":true}"#.to_owned(),
1830 r#"{"state":"in_flight","resumable":false}"#.to_owned(),
1831 ];
1832 for (record, state) in records().iter().zip(states) {
1833 let golden = format!("{head}{state}}}");
1834 assert_eq!(encode_replay_record(record).as_bytes(), golden.as_bytes());
1835 }
1836 let quota = QuotaState {
1837 window_start: 1_700_000_000_000,
1838 ops: 3,
1839 bytes: u64::MAX,
1840 };
1841 let golden =
1842 b"\x01{\"window_start\":1700000000000,\"ops\":3,\"bytes\":18446744073709551615}";
1843 assert_eq!(encode_quota_state("a).as_bytes(), golden);
1844 assert_eq!(encode_ref_id(&[4; 32]).as_bytes(), [4; 32]);
1845 assert_eq!(
1846 encode_u64(u64::MAX - 1).as_bytes(),
1847 [0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xfe]
1848 );
1849 assert_eq!(encode_u64(1).as_bytes(), [0, 0, 0, 0, 0, 0, 0, 1]);
1850 assert_eq!(encode_u32(1).as_bytes(), [0, 0, 0, 1]);
1851 }
1852
1853 #[test]
1854 fn coordinator_codecs_roundtrip_and_golden_bytes() {
1855 let namespace = NamespaceRecord {
1856 created_at_ms: 1_700_000_000_000,
1857 config_version: 1,
1858 };
1859 let repo = RepoRecord {
1860 created_at_ms: u64::MAX,
1861 };
1862 let namespace_value = encode_namespace_record(&namespace);
1863 let repo_value = encode_repo_record(&repo);
1864 assert_eq!(
1865 namespace_value.as_bytes(),
1866 b"\x01{\"created_at_ms\":1700000000000,\"config_version\":1}"
1867 );
1868 assert_eq!(
1869 repo_value.as_bytes(),
1870 b"\x01{\"created_at_ms\":18446744073709551615}"
1871 );
1872 assert_eq!(
1873 decode_namespace_record(&namespace_value).unwrap(),
1874 namespace
1875 );
1876 assert_eq!(decode_repo_record(&repo_value).unwrap(), repo);
1877 let public = RepoVisibilityV1 {
1878 visibility: StoredVisibility::Public,
1879 last_created_ms: 0,
1880 last_statement_id: None,
1881 changed_ms: None,
1882 };
1883 let private = RepoVisibilityV1 {
1884 visibility: StoredVisibility::Private,
1885 last_created_ms: 1_700_000_000_000,
1886 last_statement_id: Some("ab".repeat(32)),
1887 changed_ms: Some(1_700_000_000_001),
1888 };
1889 let public_value = encode_repo_visibility(&public);
1890 let private_value = encode_repo_visibility(&private);
1891 assert_eq!(
1892 public_value.as_bytes(),
1893 b"\x01{\"visibility\":\"public\",\"last_created_ms\":0,\"last_statement_id\":null}"
1894 );
1895 assert_eq!(
1896 private_value.as_bytes(),
1897 format!(
1898 "\x01{{\"visibility\":\"private\",\"last_created_ms\":1700000000000,\"last_statement_id\":\"{}\",\"changed_ms\":1700000000001}}",
1899 "ab".repeat(32)
1900 )
1901 .as_bytes()
1902 );
1903 assert_eq!(decode_repo_visibility(&public_value).unwrap(), public);
1904 assert_eq!(decode_repo_visibility(&private_value).unwrap(), private);
1905 let legacy = decode_repo_visibility(&Value::new(
1908 b"\x01{\"visibility\":\"private\",\"last_created_ms\":7,\"last_statement_id\":null}"
1909 .to_vec(),
1910 ))
1911 .unwrap();
1912 assert_eq!(legacy.changed_ms, None);
1913 assert_eq!(legacy.visibility_changed_ms(), 7);
1914 assert_eq!(private.visibility_changed_ms(), 1_700_000_000_001);
1915 for bytes in [
1916 &b""[..],
1917 b"\x02{\"visibility\":\"public\",\"last_created_ms\":0,\"last_statement_id\":null}",
1918 b"\x01{\"visibility\":\"internal\",\"last_created_ms\":0,\"last_statement_id\":null}",
1919 b"\x01{\"visibility\":\"public\"}",
1920 b"\x01{\"visibility\":\"public\",\"last_created_ms\":0,\"last_statement_id\":\"AB\",\"changed_ms\":0}",
1921 b"\x01{\"visibility\":\"public\",\"last_created_ms\":0}",
1922 b"\x01{\"visibility\":\"public\",\"last_created_ms\":0,\"last_statement_id\":\"ab\",\"changed_ms\":0,\"extra\":1}",
1923 ] {
1924 assert!(matches!(
1925 decode_repo_visibility(&Value::new(bytes.to_vec())),
1926 Err(StoreError::Corrupt(_))
1927 ));
1928 }
1929 for bytes in [
1930 &b""[..],
1931 b"\x02{\"created_at_ms\":0,\"config_version\":1}",
1932 b"\x01{\"created_at_ms\":-1,\"config_version\":1}",
1933 b"\x01{\"created_at_ms\":0,\"config_version\":0}",
1934 b"\x01{\"created_at_ms\":0,\"config_version\":1,\"extra\":1}",
1935 b"\x01{\"created_at_ms\":0}",
1936 ] {
1937 assert!(matches!(
1938 decode_namespace_record(&Value::new(bytes.to_vec())),
1939 Err(StoreError::Corrupt(_))
1940 ));
1941 }
1942 for bytes in [
1943 &b""[..],
1944 b"\x02{\"created_at_ms\":0}",
1945 b"\x01{\"created_at_ms\":-1}",
1946 b"\x01{\"created_at_ms\":0,\"extra\":1}",
1947 b"\x01{}",
1948 ] {
1949 assert!(matches!(
1950 decode_repo_record(&Value::new(bytes.to_vec())),
1951 Err(StoreError::Corrupt(_))
1952 ));
1953 }
1954 }
1955
1956 #[test]
1957 fn lease_codecs_roundtrip_and_golden_bytes() {
1958 let epoch = EpochLease {
1959 authority_ready: None,
1960 authority_generation: None,
1961 epoch: 7,
1962 expires_at_ms: 30000,
1963 config_version: 2,
1964 };
1965 let shard = LeasedShard {
1966 authority_generation: None,
1967 acked_authority_generation: None,
1968 epoch: 7,
1969 expires_at_ms: 30000,
1970 acked_epoch: 6,
1971 relay_watermark_ms: 123,
1972 sweep_due_ms: 30000,
1973 };
1974 let recovery = LeaseRecovery {
1975 authority_fence: None,
1976 authority_ready: None,
1977 activation_only: None,
1978 resumed_at_ms: 100_000,
1979 };
1980 let epoch_value = encode_epoch_lease(&epoch);
1981 let shard_value = encode_leased_shard(&shard);
1982 let recovery_value = encode_lease_recovery(&recovery);
1983 assert_eq!(
1984 epoch_value.as_bytes(),
1985 b"\x01{\"epoch\":7,\"expires_at_ms\":30000,\"config_version\":2}"
1986 );
1987 assert_eq!(
1988 shard_value.as_bytes(),
1989 b"\x01{\"epoch\":7,\"expires_at_ms\":30000,\"acked_epoch\":6,\"relay_watermark_ms\":123,\"sweep_due_ms\":30000}"
1990 );
1991 assert_eq!(recovery_value.as_bytes(), b"\x01{\"resumed_at_ms\":100000}");
1992 assert_eq!(decode_epoch_lease(&epoch_value).unwrap(), epoch);
1993 assert_eq!(decode_leased_shard(&shard_value).unwrap(), shard);
1994 assert_eq!(decode_lease_recovery(&recovery_value).unwrap(), recovery);
1995 for version in [0, 2, 255] {
1996 for value in [&epoch_value, &shard_value, &recovery_value] {
1997 let mut bytes = value.as_bytes().to_vec();
1998 bytes[0] = version;
1999 let bad = Value::new(bytes);
2000 assert!(decode_epoch_lease(&bad).is_err());
2001 assert!(decode_leased_shard(&bad).is_err());
2002 assert!(decode_lease_recovery(&bad).is_err());
2003 }
2004 }
2005 for body in [
2006 "{\"epoch\":7,\"expires_at_ms\":30000,\"config_version\":2,\"extra\":0}",
2007 "{\"epoch\":7,\"expires_at_ms\":30000,\"config_version\":0}",
2008 "{\"epoch\":7,\"expires_at_ms\":30000}",
2009 ] {
2010 assert!(
2011 decode_epoch_lease(&Value::new([&[CODEC_V1][..], body.as_bytes()].concat()))
2012 .is_err()
2013 );
2014 }
2015 for body in [
2016 "{\"epoch\":7,\"expires_at_ms\":30000,\"acked_epoch\":6,\"extra\":0}",
2017 "{\"epoch\":7,\"expires_at_ms\":30000}",
2018 ] {
2019 assert!(
2020 decode_leased_shard(&Value::new([&[CODEC_V1][..], body.as_bytes()].concat()))
2021 .is_err()
2022 );
2023 }
2024 assert!(
2025 decode_lease_recovery(&Value::new(
2026 b"\x01{\"resumed_at_ms\":100000,\"extra\":0}".to_vec()
2027 ))
2028 .is_err()
2029 );
2030 assert!(
2031 decode_lease_recovery(&Value::new(b"\x02{\"resumed_at_ms\":100000}".to_vec())).is_err()
2032 );
2033 assert!(
2034 decode_leased_shard(&Value::new(
2035 b"\x02{\"epoch\":7,\"expires_at_ms\":30000,\"acked_epoch\":6}".to_vec()
2036 ))
2037 .is_err()
2038 );
2039 }
2040
2041 #[test]
2042 fn content_index_codecs_golden_bytes() {
2043 let block = BlockEntry {
2044 reason: "dmca".into(),
2045 blocked_at_ms: 7,
2046 };
2047 let state = ObjectState {
2048 seq: u64::MAX,
2049 changed_at_ms: 1_700_000_000_000,
2050 holders: 2,
2051 deleting: true,
2052 };
2053 let holder = HolderRecord::new(5, [0xab; 32]);
2054 let holder_golden = format!("\x01{{\"seq\":5,\"op_id\":\"{}\"}}", "ab".repeat(32));
2055 assert_eq!(encode_holder(&holder).as_bytes(), holder_golden.as_bytes());
2056 assert_eq!(decode_holder(&encode_holder(&holder)).unwrap(), holder);
2057 for bad in [
2058 &b"\x01{\"seq\":5}"[..],
2059 b"\x01{\"seq\":5,\"op_id\":\"zz\"}",
2060 b"\x02{}",
2061 b"",
2062 ] {
2063 let bad = Value::new(bad.to_vec());
2064 assert!(matches!(decode_holder(&bad), Err(StoreError::Corrupt(_))));
2065 }
2066 let cases: [(Value, &[u8]); 3] = [
2067 (encode_hold(9), b"\x01{\"expires_at_ms\":9}"),
2068 (
2069 encode_block_entry(&block),
2070 b"\x01{\"reason\":\"dmca\",\"blocked_at_ms\":7}",
2071 ),
2072 (
2073 encode_object_state(&state),
2074 b"\x01{\"seq\":18446744073709551615,\"changed_at_ms\":1700000000000,\"holders\":2,\"deleting\":true}",
2075 ),
2076 ];
2077 for (value, golden) in &cases {
2078 assert_eq!(value.as_bytes(), *golden);
2079 }
2080 assert_eq!(decode_hold(&cases[0].0).unwrap(), 9);
2081 assert_eq!(decode_block_entry(&cases[1].0).unwrap(), block);
2082 assert_eq!(decode_object_state(&cases[2].0).unwrap(), state);
2083 for bad in [
2084 &b"\x02{\"expires_at_ms\":9}"[..],
2085 b"\x01{\"expires_at_ms\":-1}",
2086 b"\x01{\"seq\":0,\"changed_at_ms\":0}",
2087 ] {
2088 let v = Value::new(bad.to_vec());
2089 assert!(matches!(decode_hold(&v), Err(StoreError::Corrupt(_))));
2090 assert!(matches!(
2091 decode_block_entry(&v),
2092 Err(StoreError::Corrupt(_))
2093 ));
2094 assert!(matches!(
2095 decode_object_state(&v),
2096 Err(StoreError::Corrupt(_))
2097 ));
2098 }
2099 }
2100
2101 #[test]
2102 fn unknown_version_byte_is_corrupt() {
2103 let mut bytes = encode_replay_record(&records()[0]).as_bytes().to_vec();
2104 bytes[0] = 0x02;
2105 let bumped = Value::new(bytes);
2106 assert!(matches!(
2107 decode_replay_record(&bumped),
2108 Err(StoreError::Corrupt(_))
2109 ));
2110 assert!(matches!(
2111 decode_quota_state(&bumped),
2112 Err(StoreError::Corrupt(_))
2113 ));
2114 for bad in [
2115 &b""[..],
2116 b"\x01{",
2117 b"\x01{\"window_start\":0,\"ops\":0,\"bytes\":0,\"extra\":1}",
2118 b"\x01{\"fingerprint\":\"00\",\"expires_at_ms\":0,\"state\":{\"state\":\"in_flight\",\"resumable\":true}}",
2119 ] {
2120 let v = Value::new(bad.to_vec());
2121 assert!(matches!(decode_replay_record(&v), Err(StoreError::Corrupt(_))));
2122 assert!(matches!(decode_quota_state(&v), Err(StoreError::Corrupt(_))));
2123 }
2124 let retryable = format!(
2125 "\x01{{\"fingerprint\":\"{}\",\"expires_at_ms\":0,\"state\":{{\"state\":\"committed\",\"result\":{{\"kind\":\"rejected\",\"code\":\"unavailable\",\"message\":\"m\"}}}}}}",
2126 "00".repeat(32)
2127 );
2128 assert!(matches!(
2129 decode_replay_record(&Value::new(retryable.into_bytes())),
2130 Err(StoreError::Corrupt(_))
2131 ));
2132 assert!(matches!(
2133 decode_ref_id(&Value::new(vec![0; 31])),
2134 Err(StoreError::Corrupt(_))
2135 ));
2136 assert!(matches!(
2137 decode_u64(&Value::new(vec![0; 4])),
2138 Err(StoreError::Corrupt(_))
2139 ));
2140 }
2141
2142 #[test]
2143 fn namespace_quota_binary_goldens_and_corruption() {
2144 let usage = NamespaceUsage { ops: 2, bytes: 258 };
2145 assert_eq!(
2146 encode_namespace_usage(usage).as_bytes(),
2147 b"\0\0\0\0\0\0\0\x02\0\0\0\0\0\0\x01\x02"
2148 );
2149 assert_eq!(
2150 decode_namespace_usage(&encode_namespace_usage(usage)).unwrap(),
2151 usage
2152 );
2153 let view = NamespaceView {
2154 total: usage,
2155 pushed: NamespaceUsage { ops: 1, bytes: 1 },
2156 observed_at_ms: 60_000,
2157 };
2158 assert_eq!(encode_namespace_view(view).as_bytes().len(), 40);
2159 assert_eq!(
2160 decode_namespace_view(&encode_namespace_view(view)).unwrap(),
2161 view
2162 );
2163 assert!(matches!(
2164 decode_namespace_usage(&Value::new(vec![0; 15])),
2165 Err(StoreError::Corrupt(_))
2166 ));
2167 assert!(matches!(
2168 decode_namespace_view(&Value::new(vec![0; 39])),
2169 Err(StoreError::Corrupt(_))
2170 ));
2171 let invalid = NamespaceView {
2172 pushed: NamespaceUsage { ops: 3, bytes: 1 },
2173 ..view
2174 };
2175 assert!(matches!(
2176 decode_namespace_view(&encode_namespace_view(invalid)),
2177 Err(StoreError::Corrupt(_))
2178 ));
2179 }
2180}