Skip to main content

mkit_server/store/
codec.rs

1//! Value codecs. Structured values are `serde_json` behind a leading
2//! version byte ([`CODEC_V1`]); refs are the raw 32-byte id and integers
3//! are raw big-endian. Decoding an unknown version, a wrong length or a
4//! value that fails validation is [`StoreError::Corrupt`].
5
6use 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/// The namespace coordinator record. The first configuration version is 1.
28#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
29#[serde(deny_unknown_fields)]
30pub struct NamespaceRecord {
31    /// Creation time, Unix milliseconds from the business clock.
32    pub created_at_ms: u64,
33    /// Namespace configuration version, starting at 1.
34    pub config_version: u64,
35}
36
37/// The repository coordinator record.
38#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
39#[serde(deny_unknown_fields)]
40pub struct RepoRecord {
41    /// Creation time, Unix milliseconds from the business clock.
42    pub created_at_ms: u64,
43}
44
45/// A repository's stored visibility (SPEC-WRITE-GRANTS §9.1). A missing
46/// `rv` row inherits the deployment default; the row may exist before `rr` does.
47#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
48#[serde(deny_unknown_fields)]
49pub struct RepoVisibilityV1 {
50    /// The repository's visibility.
51    pub visibility: StoredVisibility,
52    /// The newer of the last accepted statement's `created` and the last
53    /// envelope write's time; 0 before any write.
54    pub last_created_ms: u64,
55    /// The last accepted statement's id, 64 lowercase hex. The key is
56    /// always present in a stored row (`null` when none): a row missing it
57    /// is corrupt.
58    #[serde(deserialize_with = "present_option")]
59    pub last_statement_id: Option<String>,
60    /// The server clock when the row last changed visibility (envelope or
61    /// statement), Unix milliseconds. URL tokens issued at or before it are
62    /// refused (SPEC-WRITE-GRANTS §9.4). Absent in a row stored before the
63    /// field existed; read it through [`Self::visibility_changed_ms`].
64    #[serde(default, skip_serializing_if = "Option::is_none")]
65    pub changed_ms: Option<u64>,
66}
67
68impl RepoVisibilityV1 {
69    /// The time of the last visibility change. A legacy row has none, so
70    /// its `last_created_ms` stands in: the envelope path stored server
71    /// time there and the statement path the accepted `created`, never
72    /// earlier than the change, so the fallback fails closed.
73    #[must_use]
74    pub fn visibility_changed_ms(&self) -> u64 {
75        self.changed_ms.unwrap_or(self.last_created_ms)
76    }
77}
78
79/// Deserialize an `Option` whose key must be present (serde would default
80/// a missing `Option` field to `None`).
81fn 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/// The visibility values a `rv` row stores.
90#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
91#[serde(rename_all = "lowercase")]
92pub enum StoredVisibility {
93    /// Readable by every caller, anonymous or signed.
94    Public,
95    /// Readable only by a signed, authorized caller (§9.3).
96    Private,
97}
98
99/// The ref shard's durable copy of its coordinator epoch lease.
100#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
101#[serde(deny_unknown_fields)]
102pub struct EpochLease {
103    /// Initial generation-zero activation barrier completed before this grant.
104    #[serde(default, skip_serializing_if = "Option::is_none")]
105    pub authority_ready: Option<bool>,
106    /// Epoch against which the shard can authorize writes.
107    pub epoch: u64,
108    /// Lease expiry, Unix milliseconds from the pipeline clock.
109    pub expires_at_ms: u64,
110    /// Namespace configuration version at grant, starting at 1.
111    pub config_version: u64,
112    /// Independent authority generation; absent in legacy unfenced codecs.
113    #[serde(default, skip_serializing_if = "Option::is_none")]
114    pub authority_generation: Option<u64>,
115}
116
117/// The coordinator's durable lease grant and installation acknowledgement.
118///
119/// `acked_epoch = n` means the shard durably holds epoch at least `n`, or
120/// every older-epoch write is already past its deadline. A live row's
121/// acknowledgement is preserved by renewal and raised only after a
122/// committed shard push. An absent or expired row may acknowledge its grant
123/// immediately because older writes are past their commit deadlines.
124#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
125#[serde(deny_unknown_fields)]
126pub struct LeasedShard {
127    /// Epoch granted by the most recent renewal.
128    pub epoch: u64,
129    /// Granted expiry; identical to the shard copy at grant.
130    pub expires_at_ms: u64,
131    /// Epoch whose installation in the shard has been acknowledged.
132    pub acked_epoch: u64,
133    /// Authority generation granted alongside the epoch.
134    #[serde(default, skip_serializing_if = "Option::is_none")]
135    pub authority_generation: Option<u64>,
136    /// Raised only after a committed push, or expiry of every older lease.
137    #[serde(default, skip_serializing_if = "Option::is_none")]
138    pub acked_authority_generation: Option<u64>,
139    /// Greatest source relay lower bound observed for this shard.
140    pub relay_watermark_ms: u64,
141    /// Due time of the one sweep timer owned by this row.
142    pub sweep_due_ms: u64,
143}
144
145/// Declared coordinator recovery, retained until a later recovery overwrites it.
146#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
147#[serde(deny_unknown_fields)]
148pub struct LeaseRecovery {
149    /// Durable namespace authority mode, independent of business creation.
150    #[serde(default, skip_serializing_if = "Option::is_none")]
151    pub authority_fence: Option<bool>,
152    /// Initial barrier is complete; fenced grants can now be installed.
153    #[serde(default, skip_serializing_if = "Option::is_none")]
154    pub authority_ready: Option<bool>,
155    /// This row carries activation state without declaring a real recovery.
156    #[serde(default, skip_serializing_if = "Option::is_none")]
157    pub activation_only: Option<bool>,
158    /// Recovery time, Unix milliseconds from the pipeline clock.
159    pub resumed_at_ms: u64,
160}
161
162impl LeaseRecovery {
163    /// Actual recovery time; a pure activation marker creates no holdoff.
164    #[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/// The last consistent Worker snapshot of one partition. A zero export time
171/// and empty key mark a timer that has been seeded but has not fired yet.
172#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
173#[serde(deny_unknown_fields)]
174pub struct BackupStateV1 {
175    /// Export time encoded in the last uploaded snapshot, Unix milliseconds.
176    pub last_export_ms: u64,
177    /// BLAKE3 of the last uploaded portable export.
178    pub digest: Hash,
179    /// R2 object key of the last upload.
180    pub r2_key: String,
181    /// Time of the last successful upload, Unix milliseconds.
182    pub last_upload_ms: u64,
183}
184
185/// An open upload ticket. Its audience is bound by the shard's deployment.
186#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
187#[serde(deny_unknown_fields)]
188pub struct TicketV1 {
189    /// Authority generation authorized when this ticket was created.
190    #[serde(default, skip_serializing_if = "Option::is_none")]
191    pub authority_generation: Option<u64>,
192    /// Repository name within the partition's namespace.
193    #[serde(with = "repo_json")]
194    pub repo: RepoName,
195    /// Ref authorized to consume the ticket.
196    pub ref_name: String,
197    /// Authorized signer.
198    #[serde(with = "hash_json")]
199    pub signer: Hash,
200    /// Expected pack id.
201    #[serde(with = "hash_json")]
202    pub pack_id: Hash,
203    /// Declared nonzero pack size.
204    pub bytes: u64,
205    /// Power-of-two part size, at least the protocol minimum.
206    pub part_size: u64,
207    /// Expiry, Unix milliseconds.
208    pub expires_at_ms: u64,
209    /// Creation, Unix milliseconds.
210    pub created_at_ms: u64,
211    /// One durable outcome id, including synthetic ids for default admission.
212    pub reservation_id: String,
213    /// Backend multipart upload session, when allocated.
214    #[serde(with = "optional_bytes_hex_json")]
215    pub upload_session: Option<Vec<u8>>,
216}
217
218/// The hooks protocol's terminal abort reasons.
219#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
220#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
221pub enum AbortReason {
222    /// A policy or pre-apply refusal without a more specific reason.
223    Unspecified,
224    /// A guarded ref update failed.
225    RefConflict,
226    /// An epoch changed.
227    EpochMismatch,
228    /// A required pack is missing.
229    PackMissing,
230    /// A concurrent request won the replay race.
231    ReplayRace,
232    /// An internal failure prevented the apply.
233    Internal,
234    /// Reconcile found an abandoned pending reservation.
235    Abandoned,
236}
237
238/// The operation a pending reservation will settle.
239#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
240#[serde(rename_all = "snake_case")]
241pub enum PendingOp {
242    /// A directly admitted write or `BeginUpload`.
243    Write,
244    /// An admitted HTTP read.
245    Read,
246}
247
248/// A ref changed by a committed reservation.
249#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
250#[serde(deny_unknown_fields)]
251pub struct OutcomeRef {
252    /// Full ref name.
253    pub name: String,
254    /// New ref target; absent on deletion.
255    #[serde(with = "optional_hash_json")]
256    pub new: Option<Hash>,
257    /// Whether this ref was deleted.
258    pub deleted: bool,
259}
260
261/// The one durable reservation arbiter, replaced under an Equals guard.
262///
263/// Unknown state tags fail decoding, so older readers fail closed.
264#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
265#[serde(tag = "state", rename_all = "snake_case", deny_unknown_fields)]
266pub enum ReservationV1 {
267    /// Admission's durable pre-apply arbiter.
268    Pending {
269        /// Full wire repository identity (or bare single-deployment name).
270        repository: String,
271        /// Creation time, Unix milliseconds.
272        created_at_ms: u64,
273        /// Earliest safe reconciliation time, Unix milliseconds.
274        reconcile_at_ms: u64,
275        /// Write or read semantics.
276        op: PendingOp,
277    },
278    /// Successful `BeginUpload`, awaiting ticket consumption or expiry.
279    Ticketed {
280        /// Bound ticket id.
281        #[serde(with = "hash_json")]
282        ticket_id: Hash,
283    },
284    /// A committed apply, including its byte accounting and ref changes.
285    Committed {
286        /// Full wire repository identity (or a bare single-deployment name).
287        repository: String,
288        /// Outcome time, Unix milliseconds.
289        occurred_at_ms: u64,
290        /// Stored bytes.
291        bytes_stored: u64,
292        /// Bytes new to the repository.
293        new_to_repo: u64,
294        /// Bytes new to the store.
295        new_to_store: u64,
296        /// Ref changes included in this apply.
297        refs: Vec<OutcomeRef>,
298    },
299    /// Failed apply or abandoned pending reservation.
300    Aborted {
301        /// Full wire repository identity (or a bare single-deployment name).
302        repository: String,
303        /// Outcome time, Unix milliseconds.
304        occurred_at_ms: u64,
305        /// Stable hooks reason.
306        reason: AbortReason,
307        /// Safe diagnostic detail, bounded to 512 UTF-8 bytes.
308        detail: String,
309    },
310    /// An unconsumed ticket expired.
311    Expired {
312        /// Full wire repository identity (or a bare single-deployment name).
313        repository: String,
314        /// Outcome time, Unix milliseconds.
315        occurred_at_ms: u64,
316    },
317    /// An admitted HTTP read, including partial delivery.
318    ReadServed {
319        /// Full wire repository identity (or bare single-deployment name).
320        repository: String,
321        /// Completion time, Unix milliseconds.
322        occurred_at_ms: u64,
323        /// Object identifier served.
324        #[serde(with = "hash_json")]
325        object: Hash,
326        /// Actual body bytes sent; zero for HEAD.
327        bytes_served: u64,
328    },
329}
330
331/// An idempotent relay of upserts and deletes to one partition.
332#[derive(Debug, Clone, PartialEq, Eq)]
333pub struct RelayV1 {
334    /// Writer plan-time lower bound on commit time. V1 changed in place before deployment.
335    pub at_ms: u64,
336    /// Destination partition.
337    pub target: Partition,
338    /// Idempotent key/value upserts.
339    pub puts: Vec<(Key, Value)>,
340    /// Keys removed after this row's upserts.
341    pub deletes: Vec<Key>,
342}
343
344/// Maximum retained targets in a persistent relay scan cycle.
345pub const MAX_BLOCKED_TARGETS: usize = 32;
346
347/// Source-local scan progress. Every retained row at or below `cursor`
348/// belongs to a blocked target; new rows beyond `cycle_end` wait for the
349/// next cycle. The relay checkpoints this row under a guard on its prior value,
350/// atomically deleting delivered queue rows. Timer rescheduling is separate.
351#[derive(Debug, Clone, PartialEq, Eq)]
352pub struct RelayScanV1 {
353    /// Source `os` observed when the scan cycle started.
354    pub cycle_end: u64,
355    /// Last inspected relay sequence, or zero before the first row.
356    pub cursor: u64,
357    /// Sorted, unique retained targets in [`Partition`] order.
358    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/// Terminal outcome backlog. Relay rows are excluded.
370#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
371#[serde(deny_unknown_fields)]
372pub struct Backlog {
373    /// Number of terminal outcome rows.
374    pub rows: u64,
375    /// Sum of terminal outcome key lengths plus encoded value lengths.
376    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    // serde with-module serialization requires a reference to the field.
408    #[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
489/// Version byte of every structured value this binary writes.
490pub 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
570/// Every [`Code`], to invert [`Code::as_str`].
571const 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/// Encode a namespace coordinator record.
615#[must_use]
616pub fn encode_namespace_record(record: &NamespaceRecord) -> Value {
617    encode_json(record)
618}
619
620/// Decode a namespace coordinator record. Configuration version zero is
621/// invalid: the namespace's first version is 1.
622pub 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/// Encode a repository coordinator record.
631#[must_use]
632pub fn encode_repo_record(record: &RepoRecord) -> Value {
633    encode_json(record)
634}
635
636/// Decode a repository coordinator record.
637pub fn decode_repo_record(value: &Value) -> Result<RepoRecord, StoreError> {
638    decode_json(value, "bad repo record")
639}
640
641/// Encode a repository visibility row.
642#[must_use]
643pub fn encode_repo_visibility(row: &RepoVisibilityV1) -> Value {
644    encode_json(row)
645}
646
647/// Decode a repository visibility row. The statement id is canonical: 64
648/// lowercase hexadecimal digits.
649pub 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/// Encode a ref shard epoch lease.
660#[must_use]
661pub fn encode_epoch_lease(lease: &EpochLease) -> Value {
662    encode_json(lease)
663}
664
665/// Decode an epoch lease; namespace configuration versions start at 1.
666pub 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/// Encode a coordinator lease-table row.
678#[must_use]
679pub fn encode_leased_shard(lease: &LeasedShard) -> Value {
680    encode_json(lease)
681}
682
683/// Decode a coordinator lease-table row.
684pub fn decode_leased_shard(value: &Value) -> Result<LeasedShard, StoreError> {
685    decode_json(value, "bad leased shard")
686}
687
688/// Encode a declared lease-table recovery marker.
689#[must_use]
690pub fn encode_lease_recovery(recovery: &LeaseRecovery) -> Value {
691    encode_json(recovery)
692}
693
694/// Decode a declared lease-table recovery marker.
695pub 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/// Encode a per-partition Worker backup state.
708#[must_use]
709pub fn encode_backup_state(state: &BackupStateV1) -> Value {
710    encode_json(state)
711}
712
713/// Decode a per-partition Worker backup state.
714pub fn decode_backup_state(value: &Value) -> Result<BackupStateV1, StoreError> {
715    decode_json(value, "bad backup state")
716}
717
718/// Validate ticket semantics before opening or after decoding a row.
719pub 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/// Encode an upload ticket.
734#[must_use]
735pub fn encode_ticket(ticket: &TicketV1) -> Value {
736    encode_json(ticket)
737}
738
739/// Decode an upload ticket, validating ref, reservation, geometry and lifetime.
740pub 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/// Encode a pending, ticket-backed, or terminal reservation.
748#[must_use]
749pub fn encode_reservation(reservation: &ReservationV1) -> Value {
750    encode_json(reservation)
751}
752
753/// Decode a reservation; unknown states, identities and malformed outcomes fail closed.
754pub 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
821/// Encode idempotent relay operations. A malformed target returns Invalid.
822pub 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
858/// Decode a relay target and its bounded idempotent operations.
859pub 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
896/// Encode a repository object-index value. All integers are big-endian.
897/// Layout: version, frame offset/length (u64 each), wire type (u8),
898/// decoded size (u64), chain depth (u32), base-present (u8), optional base id.
899pub 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
918/// Decode and validate a repository object-index value.
919pub 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
974/// Encode bounded, canonical source relay scan progress.
975pub 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
997/// Decode source relay scan progress, rejecting malformed or unbounded state.
998pub 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/// Encode the terminal outcome backlog.
1021#[must_use]
1022pub fn encode_backlog(backlog: &Backlog) -> Value {
1023    encode_json(backlog)
1024}
1025
1026/// Decode the terminal outcome backlog.
1027pub 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/// Encode a replay record.
1037#[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
1088/// Decode a replay record.
1089pub 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/// Encode a quota window's usage.
1151#[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
1160/// Decode a quota window's usage.
1161pub 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/// Two big-endian u64 counters, ops then bytes. Used by qs, qc and qt.
1171#[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
1176/// Decode a fixed-width namespace counter, rejecting malformed rows.
1177pub 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/// Three usage fields and a read timestamp, each big-endian u64.
1192#[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
1206/// Decode a view and reject a contribution larger than the aggregate.
1207pub 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/// Encode a `ContentIndex` GC hold: its expiry, Unix ms.
1238#[must_use]
1239pub fn encode_hold(expires_at_ms: u64) -> Value {
1240    encode_json(&HoldV1 { expires_at_ms })
1241}
1242
1243/// Decode a `ContentIndex` GC hold's expiry.
1244pub 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/// Encode a `ContentIndex` holder row (R-131).
1250#[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
1258/// Decode a `ContentIndex` holder row.
1259pub 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/// Encode a blocklist entry.
1265#[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
1273/// Decode a blocklist entry.
1274pub 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/// Encode a `ContentIndex` object state.
1283#[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
1293/// Decode a `ContentIndex` object state.
1294pub 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/// A ref value: the raw 32-byte id.
1305#[must_use]
1306pub fn encode_ref_id(id: &Hash) -> Value {
1307    Value::new(id.to_vec())
1308}
1309
1310/// Decode a ref value.
1311pub 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/// A be64 integer (the grant epoch).
1316#[must_use]
1317pub fn encode_u64(n: u64) -> Value {
1318    Value::new(n.to_be_bytes().to_vec())
1319}
1320
1321/// Decode a be64 integer.
1322pub 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/// A be32 integer (the layout version).
1329#[must_use]
1330pub fn encode_u32(n: u32) -> Value {
1331    Value::new(n.to_be_bytes().to_vec())
1332}
1333
1334/// Decode a be32 integer.
1335pub 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(&quota)).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(&quota).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        // A legacy 3-key row has no change time and falls back to its
1906        // last accepted creation time.
1907        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}