Skip to main content

mkit_server/store/
tickets.rs

1//! Pure ticket/counter and ref-shard membership fragments. Callers supply
2//! time and read snapshots, compose these fragments with their ref/replay
3//! writes, and apply the complete batch in one partition.
4
5use std::collections::BTreeSet;
6
7use super::codec::{self, TicketV1};
8use super::outbox::{OutboxBuilder, guard};
9use super::{BlobKey, Key, Partition, Precondition, StoreError, Value, Write, keys as layout};
10use crate::pipeline::ShardMap;
11use crate::repo::{RepoId, RepoName};
12use crate::timers::registry::kinds;
13use mkit_core::hash::Hash;
14
15/// Caller-validated upload geometry and binding, plus the business clock.
16#[derive(Debug, Clone)]
17pub struct TicketSpec {
18    /// Authority generation authorized when this ticket was created.
19    pub authority_generation: Option<u64>,
20    /// Repository name within this partition's namespace.
21    pub repo: RepoName,
22    /// Target ref (wire normalization is WP-1.9).
23    pub ref_name: String,
24    /// Authenticated signer.
25    pub signer: Hash,
26    /// Pack commitment.
27    pub pack_id: Hash,
28    /// Positive committed byte count.
29    pub bytes: u64,
30    /// Power-of-two part geometry.
31    pub part_size: u64,
32    /// Expiry, strictly less than seven days after creation.
33    pub expires_at_ms: u64,
34    /// Creation time.
35    pub created_at_ms: u64,
36    /// Business-clock observation used to classify an existing ticket.
37    pub now_ms: u64,
38    /// Admission id, or the synthetic replay-scope id.
39    pub reservation_id: String,
40    /// Optional backend multipart session identifier.
41    pub upload_session: Option<Vec<u8>>,
42}
43
44impl TicketSpec {
45    fn record(&self) -> TicketV1 {
46        TicketV1 {
47            authority_generation: self.authority_generation,
48            repo: self.repo.clone(),
49            ref_name: self.ref_name.clone(),
50            signer: self.signer,
51            pack_id: self.pack_id,
52            bytes: self.bytes,
53            part_size: self.part_size,
54            expires_at_ms: self.expires_at_ms,
55            created_at_ms: self.created_at_ms,
56            reservation_id: self.reservation_id.clone(),
57            upload_session: self.upload_session.clone(),
58        }
59    }
60}
61
62/// Open-ticket caps, decided by the embedding RPC.
63#[derive(Debug, Clone, Copy)]
64pub struct TicketCaps {
65    /// Maximum across all signers of the ref.
66    pub per_ref: u64,
67    /// Maximum for one ref/signer pair.
68    pub per_signer: u64,
69}
70
71/// Initial read set. After reading index, fetch its ticket separately if
72/// it differs from ticket (the reservation-derived new id).
73#[derive(Debug, Clone)]
74pub struct TicketReadKeys {
75    /// Proposed ticket row.
76    pub ticket: Key,
77    /// Idempotency index.
78    pub index: Key,
79    /// Shared ref counter.
80    pub per_ref: Key,
81    /// Signer counter.
82    pub per_signer: Key,
83    /// Reservation arbiter.
84    pub reservation: Key,
85}
86
87/// Compute keys for a valid spec.
88///
89/// # Panics
90/// If the caller supplies an invalid ref/reservation id. The open planner
91/// validates untrusted specs and returns an error instead.
92#[must_use]
93pub fn keys(spec: &TicketSpec) -> TicketReadKeys {
94    TicketReadKeys {
95        ticket: layout::ticket(&ticket_id(&spec.reservation_id)),
96        index: layout::ticket_index(&spec.repo, &spec.ref_name, &spec.pack_id, &spec.signer)
97            .expect("validated ticket binding"),
98        per_ref: layout::tickets_per_ref(&spec.repo, &spec.ref_name).expect("validated ref"),
99        per_signer: layout::tickets_per_signer(&spec.repo, &spec.ref_name, &spec.signer)
100            .expect("validated ref"),
101        reservation: layout::reservation(&spec.reservation_id).expect("validated reservation id"),
102    }
103}
104
105/// Snapshot for keys, plus the ticket named by the index (None means the
106/// indexed row was read and is absent). If index names the proposed id,
107/// ticket is also sufficient; an explicitly supplied `indexed_ticket` wins.
108#[derive(Debug, Clone, Default)]
109pub struct TicketReads {
110    /// Proposed ticket value.
111    pub ticket: Option<Value>,
112    /// Ticket value fetched using the id in index.
113    pub indexed_ticket: Option<Value>,
114    /// Raw 32-byte ticket id in the idempotency index.
115    pub index: Option<Value>,
116    /// Ref counter.
117    pub per_ref: Option<Value>,
118    /// Ref/signer counter.
119    pub per_signer: Option<Value>,
120    /// Reservation row; a duplicate is rejected by Absent at apply.
121    pub reservation: Option<Value>,
122}
123
124/// Open planning leaves output vectors untouched on every error.
125#[derive(Debug)]
126pub enum TicketPlanError {
127    /// Idempotent re-creation returns the still-live ticket.
128    Existing(TicketV1),
129    /// Ref cap if true; otherwise signer cap.
130    CapExceeded { per_ref: bool },
131    /// Invalid stored data.
132    Corrupt(StoreError),
133    /// Invalid caller input.
134    Invalid(&'static str),
135}
136
137/// Deterministic, audience-local ticket identifier.
138#[must_use]
139pub fn ticket_id(reservation_id: &str) -> Hash {
140    mkit_core::hash::hash(&[b"mkit.ticket.v1\n".as_slice(), reservation_id.as_bytes()].concat())
141}
142
143fn counter(value: Option<&Value>) -> Result<u64, StoreError> {
144    let n = value.map(codec::decode_u64).transpose()?.unwrap_or(0);
145    if value.is_some() && n == 0 {
146        return Err(StoreError::Corrupt("open counter stored as zero".into()));
147    }
148    Ok(n)
149}
150
151/// Coalesce repeated edits of a shared counter without duplicating its
152/// snapshot guard. This is required when an advance consumes many tickets.
153fn adjust_counter(
154    key: Key,
155    prior: Option<&Value>,
156    increment: bool,
157    pre: &mut Vec<Precondition>,
158    writes: &mut Vec<Write>,
159) -> Result<(), StoreError> {
160    let observed = counter(prior)?;
161    let expected = guard(key.clone(), prior);
162    let existing = pre.iter().find(|p| match p {
163        Precondition::Equals(k, _) | Precondition::Absent(k) | Precondition::Present(k) => {
164            k == &key
165        }
166        Precondition::NotAfter(_) => false,
167    });
168    if existing.is_some_and(|p| p != &expected) {
169        return Err(StoreError::Invalid("inconsistent counter snapshots".into()));
170    }
171    let position = writes.iter().rposition(|w| match w {
172        Write::Put(k, _) | Write::Delete(k) => k == &key,
173    });
174    let current = match position.map(|i| &writes[i]) {
175        Some(Write::Put(_, value)) => counter(Some(value))?,
176        Some(Write::Delete(_)) => 0,
177        None => observed,
178    };
179    let next = if increment {
180        current.checked_add(1)
181    } else {
182        current.checked_sub(1)
183    }
184    .ok_or_else(|| StoreError::Corrupt("open counter overflow/underflow".into()))?;
185    if existing.is_none() {
186        pre.push(expected);
187    }
188    let write = if next == 0 {
189        Write::Delete(key)
190    } else {
191        Write::Put(key, codec::encode_u64(next))
192    };
193    if let Some(i) = position {
194        writes[i] = write;
195    } else {
196        writes.push(write);
197    }
198    Ok(())
199}
200
201/// Open a ticket and its Ticketed reservation in one fragment.
202/// Only one open per ref may be composed in a batch, avoiding stale cap
203/// snapshots. The reservation uses no sequence/backlog, so its local
204/// builder requires neither os nor oc reads; callers may also finish their
205/// own builder.
206// The fixed public contract returns Existing(TicketV1) by value.
207#[allow(clippy::result_large_err, clippy::too_many_lines)] // One atomic ticket, cap, index and reservation fragment.
208pub fn plan_ticket_open(
209    spec: &TicketSpec,
210    reads: &TicketReads,
211    caps: TicketCaps,
212    pre: &mut Vec<Precondition>,
213    writes: &mut Vec<Write>,
214) -> Result<Hash, TicketPlanError> {
215    let ticket = spec.record();
216    let value = codec::encode_ticket(&ticket);
217    codec::decode_ticket(&value)
218        .map_err(|_| TicketPlanError::Invalid("invalid ticket specification"))?;
219    if spec.expires_at_ms <= spec.now_ms {
220        return Err(TicketPlanError::Invalid("new ticket is already expired"));
221    }
222    let read_keys = keys(spec);
223    let id = ticket_id(&spec.reservation_id);
224    if let Some(value) = &reads.reservation
225        && !matches!(
226            codec::decode_reservation(value).map_err(TicketPlanError::Corrupt)?,
227            codec::ReservationV1::Pending {
228                op: codec::PendingOp::Write,
229                ..
230            }
231        )
232    {
233        return Err(TicketPlanError::Invalid("reservation id already in use"));
234    }
235    // A ticket row the index doesn't name as live must not exist.
236    let mut unread_indexed = None;
237    if let Some(index) = &reads.index {
238        let indexed_id = codec::decode_ref_id(index).map_err(TicketPlanError::Corrupt)?;
239        let raw = reads.indexed_ticket.as_ref().or_else(|| {
240            (indexed_id == id)
241                .then_some(reads.ticket.as_ref())
242                .flatten()
243        });
244        if let Some(raw) = raw {
245            let existing = codec::decode_ticket(raw).map_err(TicketPlanError::Corrupt)?;
246            if existing.repo != spec.repo
247                || existing.ref_name != spec.ref_name
248                || existing.signer != spec.signer
249                || existing.pack_id != spec.pack_id
250                || ticket_id(&existing.reservation_id) != indexed_id
251            {
252                return Err(TicketPlanError::Corrupt(StoreError::Corrupt(
253                    "ticket index binding mismatch".into(),
254                )));
255            }
256            if existing.expires_at_ms > spec.now_ms {
257                return Err(TicketPlanError::Existing(existing));
258            }
259        } else if indexed_id != id {
260            // The caller didn't read the indexed ticket: require it gone, so
261            // a live ticket can never be shadowed by a second one.
262            unread_indexed = Some(layout::ticket(&indexed_id));
263        }
264    }
265    if reads.ticket.is_some() {
266        return Err(TicketPlanError::Invalid("ticket id already in use"));
267    }
268    // BeginUpload creates exactly one ticket. Reject a second open for
269    // this ref in a composed batch rather than using stale cap snapshots.
270    if writes
271        .iter()
272        .any(|w| matches!(w, Write::Put(k, _) | Write::Delete(k) if k == &read_keys.per_ref))
273    {
274        return Err(TicketPlanError::Invalid(
275            "one ticket open per ref per batch",
276        ));
277    }
278    let per_ref = counter(reads.per_ref.as_ref()).map_err(TicketPlanError::Corrupt)?;
279    let per_signer = counter(reads.per_signer.as_ref()).map_err(TicketPlanError::Corrupt)?;
280    if per_ref >= caps.per_ref {
281        return Err(TicketPlanError::CapExceeded { per_ref: true });
282    }
283    if per_signer >= caps.per_signer {
284        return Err(TicketPlanError::CapExceeded { per_ref: false });
285    }
286    let (mut staged_pre, mut staged_writes) = (pre.clone(), writes.clone());
287    staged_pre.extend([
288        Precondition::Absent(read_keys.ticket.clone()),
289        guard(read_keys.index.clone(), reads.index.as_ref()),
290    ]);
291    if let Some(key) = unread_indexed {
292        staged_pre.push(Precondition::Absent(key));
293    }
294    staged_writes.extend([
295        Write::Put(read_keys.ticket, value),
296        Write::Put(read_keys.index, codec::encode_ref_id(&id)),
297    ]);
298    adjust_counter(
299        read_keys.per_ref,
300        reads.per_ref.as_ref(),
301        true,
302        &mut staged_pre,
303        &mut staged_writes,
304    )
305    .map_err(TicketPlanError::Corrupt)?;
306    adjust_counter(
307        read_keys.per_signer,
308        reads.per_signer.as_ref(),
309        true,
310        &mut staged_pre,
311        &mut staged_writes,
312    )
313    .map_err(TicketPlanError::Corrupt)?;
314    staged_writes.push(Write::Put(
315        layout::timer(spec.expires_at_ms, kinds::TICKET_EXPIRY.get(), &id),
316        Value::default(),
317    ));
318    let mut outbox = OutboxBuilder::new(None, None).map_err(TicketPlanError::Corrupt)?;
319    outbox.reserve(&spec.reservation_id, id, reads.reservation.as_ref());
320    outbox
321        .try_finish(&mut staged_pre, &mut staged_writes)
322        .map_err(TicketPlanError::Corrupt)?;
323    *pre = staged_pre;
324    *writes = staged_writes;
325    Ok(id)
326}
327
328/// Whether the caller consumed the ticket or processed its expiry.
329#[derive(Debug, Clone, Copy, PartialEq, Eq)]
330pub enum CloseReason {
331    /// Keep the timer for the expiry handler's missing-ticket no-op.
332    Consumed,
333    /// A lost pack is terminal; keep the timer for the same missing-ticket no-op.
334    Aborted,
335    /// Remove the timer along with the ticket.
336    Expired,
337    /// The timer core removes the fired expiry row with the same batch.
338    ExpiryTimerFired,
339}
340
341/// Close only the exact ticket observed (R-04). The caller writes its
342/// Committed/Expired outcome using the same batch. Errors change nothing.
343#[allow(clippy::too_many_arguments)]
344pub fn plan_ticket_close(
345    ticket_id: &Hash,
346    ticket: &TicketV1,
347    ticket_value: &Value,
348    ti_value: Option<&Value>,
349    tc: Option<&Value>,
350    tu: Option<&Value>,
351    why: CloseReason,
352    pre: &mut Vec<Precondition>,
353    writes: &mut Vec<Write>,
354) -> Result<(), StoreError> {
355    if codec::decode_ticket(ticket_value)? != *ticket
356        || self::ticket_id(&ticket.reservation_id) != *ticket_id
357    {
358        return Err(StoreError::Corrupt("ticket value/id mismatch".into()));
359    }
360    let index = layout::ticket_index(
361        &ticket.repo,
362        &ticket.ref_name,
363        &ticket.pack_id,
364        &ticket.signer,
365    )?;
366    let indexed = ti_value.map(codec::decode_ref_id).transpose()?;
367    let (mut staged_pre, mut staged_writes) = (pre.clone(), writes.clone());
368    let key = layout::ticket(ticket_id);
369    if writes
370        .iter()
371        .any(|w| matches!(w, Write::Delete(k) | Write::Put(k, _) if k == &key))
372    {
373        return Err(StoreError::Invalid(
374            "ticket already planned in batch".into(),
375        ));
376    }
377    staged_pre.push(Precondition::Equals(key.clone(), ticket_value.clone()));
378    staged_writes.push(Write::Delete(key));
379    if indexed.as_ref() == Some(ticket_id) {
380        staged_pre.push(guard(index.clone(), ti_value));
381        staged_writes.push(Write::Delete(index));
382    }
383    adjust_counter(
384        layout::tickets_per_ref(&ticket.repo, &ticket.ref_name)?,
385        tc,
386        false,
387        &mut staged_pre,
388        &mut staged_writes,
389    )?;
390    adjust_counter(
391        layout::tickets_per_signer(&ticket.repo, &ticket.ref_name, &ticket.signer)?,
392        tu,
393        false,
394        &mut staged_pre,
395        &mut staged_writes,
396    )?;
397    if why == CloseReason::Expired {
398        staged_writes.push(Write::Delete(layout::timer(
399            ticket.expires_at_ms,
400            kinds::TICKET_EXPIRY.get(),
401            ticket_id,
402        )));
403    }
404    *pre = staged_pre;
405    *writes = staged_writes;
406    Ok(())
407}
408
409/// Add immediate local memberships and queue live/published upserts for every distinct
410/// remote target. `SinglePartition` needs no relay row or sequence update.
411pub fn plan_membership(
412    repo: &RepoName,
413    packs: &[Hash],
414    source: &Partition,
415    shards: &dyn ShardMap,
416    repo_id: &RepoId,
417    outbox: &mut OutboxBuilder,
418    writes: &mut Vec<Write>,
419) {
420    debug_assert_eq!(repo, &repo_id.name, "membership repo must match its RepoId");
421    for pack in packs.iter().collect::<BTreeSet<_>>() {
422        let key = layout::membership(repo, pack);
423        let put = Write::Put(key.clone(), Value::default());
424        if !writes.contains(&put) {
425            writes.push(put);
426        }
427        let target = shards.membership(repo_id, &BlobKey::pack(*pack));
428        if target != *source {
429            outbox.relay(
430                &target,
431                vec![
432                    (key, Value::default()),
433                    (layout::published_member(repo, pack), Value::default()),
434                ],
435            );
436        }
437    }
438}
439
440#[cfg(test)]
441#[path = "tickets_tests.rs"]
442mod tests;