1use 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#[derive(Debug, Clone)]
17pub struct TicketSpec {
18 pub authority_generation: Option<u64>,
20 pub repo: RepoName,
22 pub ref_name: String,
24 pub signer: Hash,
26 pub pack_id: Hash,
28 pub bytes: u64,
30 pub part_size: u64,
32 pub expires_at_ms: u64,
34 pub created_at_ms: u64,
36 pub now_ms: u64,
38 pub reservation_id: String,
40 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#[derive(Debug, Clone, Copy)]
64pub struct TicketCaps {
65 pub per_ref: u64,
67 pub per_signer: u64,
69}
70
71#[derive(Debug, Clone)]
74pub struct TicketReadKeys {
75 pub ticket: Key,
77 pub index: Key,
79 pub per_ref: Key,
81 pub per_signer: Key,
83 pub reservation: Key,
85}
86
87#[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#[derive(Debug, Clone, Default)]
109pub struct TicketReads {
110 pub ticket: Option<Value>,
112 pub indexed_ticket: Option<Value>,
114 pub index: Option<Value>,
116 pub per_ref: Option<Value>,
118 pub per_signer: Option<Value>,
120 pub reservation: Option<Value>,
122}
123
124#[derive(Debug)]
126pub enum TicketPlanError {
127 Existing(TicketV1),
129 CapExceeded { per_ref: bool },
131 Corrupt(StoreError),
133 Invalid(&'static str),
135}
136
137#[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
151fn 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#[allow(clippy::result_large_err, clippy::too_many_lines)] pub 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 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 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 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
330pub enum CloseReason {
331 Consumed,
333 Aborted,
335 Expired,
337 ExpiryTimerFired,
339}
340
341#[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
409pub 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;