1use std::sync::atomic::{AtomicU64, Ordering};
4
5use mkit_core::hash::Hash;
6use mkit_core::repo_identity::RepositoryIdentity;
7
8use super::{DueTimer, Fired, TimerCtx, TimerHandler, TimerKind, registry::kinds};
9use crate::repo::NamespaceKey;
10use crate::rt::BoxFuture;
11use crate::store::codec::{self, ReservationV1, TicketV1};
12use crate::store::outbox::{OutboxBuilder, Terminal};
13use crate::store::tickets::{self, CloseReason};
14use crate::store::{
15 Batch, BlobKey, MultipartBlobStore, NamespaceStore, Partition, Precondition, StoreError, Write,
16 keys,
17};
18
19static ABORT_FAILURES: AtomicU64 = AtomicU64::new(0);
20
21#[must_use]
23pub fn abort_failures() -> u64 {
24 ABORT_FAILURES.load(Ordering::Relaxed)
25}
26
27#[derive(Debug)]
29pub struct TicketExpiry<B> {
30 pub blobs: B,
32}
33
34fn repository(partition: &Partition, ticket: &TicketV1) -> Result<String, StoreError> {
35 let ns = match partition {
36 Partition::Namespace(ns) => ns,
37 Partition::Ref {
38 ns,
39 repo,
40 shard_ref,
41 } if repo == &ticket.repo && shard_ref == &ticket.ref_name => ns,
42 _ => return Err(StoreError::Corrupt("ticket in wrong partition".into())),
43 };
44 let name = if ns == &NamespaceKey::deployment_default() {
45 ticket.repo.as_str().to_owned()
46 } else {
47 format!("{}/{}", ns.as_str(), ticket.repo.as_str())
48 };
49 RepositoryIdentity::parse_bare_allowed(&name)
50 .map_err(|_| StoreError::Corrupt("invalid ticket repository".into()))?;
51 Ok(name)
52}
53
54async fn verification_cleanup<S: NamespaceStore>(
58 ctx: &TimerCtx<'_, S>,
59 ticket: &TicketV1,
60 batch: &mut Batch,
61) -> Result<(), StoreError> {
62 let vs_key = keys::verification(&ticket.repo, &ticket.pack_id);
63 let rows = ctx
64 .store
65 .get_many(
66 ctx.partition,
67 &[
68 keys::membership(&ticket.repo, &ticket.pack_id),
69 vs_key.clone(),
70 keys::verify_job(&ticket.repo, &ticket.pack_id),
71 ],
72 )
73 .await?;
74 let [member, state, job] = rows.as_slice() else {
75 return Err(StoreError::Corrupt(
76 "short verification cleanup read".into(),
77 ));
78 };
79 if let (None, Some(raw)) = (member, state) {
80 batch
81 .preconditions
82 .push(Precondition::Equals(vs_key.clone(), raw.clone()));
83 batch
84 .preconditions
85 .push(Precondition::Absent(keys::membership(
86 &ticket.repo,
87 &ticket.pack_id,
88 )));
89 batch.writes.push(Write::Delete(vs_key));
90 }
91 if job.is_some() {
92 batch.writes.push(Write::Put(
93 keys::timer(
94 ctx.now_ms,
95 kinds::VERIFY.get(),
96 &crate::indexed::checkpoint::timer_reference(&ticket.repo, &ticket.pack_id),
97 ),
98 crate::store::Value::default(),
99 ));
100 }
101 Ok(())
102}
103
104impl<S: NamespaceStore, B: MultipartBlobStore> TimerHandler<S> for TicketExpiry<B> {
105 fn kind(&self) -> TimerKind {
106 kinds::TICKET_EXPIRY
107 }
108
109 fn max_per_tick(&self) -> Option<u32> {
110 Some(8)
111 }
112
113 fn fire<'a>(
114 &'a self,
115 ctx: &'a TimerCtx<'a, S>,
116 timer: &'a DueTimer,
117 ) -> BoxFuture<'a, Result<Fired, StoreError>> {
118 Self::fire_with(&self.blobs, ctx, timer)
119 }
120}
121
122impl<B: MultipartBlobStore> TicketExpiry<B> {
123 fn fire_with<'a, S: NamespaceStore>(
124 blobs: &'a B,
125 ctx: &'a TimerCtx<'a, S>,
126 timer: &'a DueTimer,
127 ) -> BoxFuture<'a, Result<Fired, StoreError>> {
128 Box::pin(async move {
129 let fire = async {
130 let id: Hash = timer.reference.as_ref().try_into().map_err(|_| {
131 StoreError::Corrupt("ticket expiry reference is not a ticket id".into())
132 })?;
133 let key = keys::ticket(&id);
134 let Some(raw) = ctx.store.get(ctx.partition, &key).await? else {
135 return Ok(Fired::Done(Batch::new()));
136 };
137 let ticket = codec::decode_ticket(&raw)?;
138 if ticket.expires_at_ms > ctx.now_ms {
139 return Ok(Fired::Reschedule {
140 due_at_ms: ticket.expires_at_ms,
141 value: timer.value.clone(),
142 batch: Batch::new(),
143 });
144 }
145
146 let rid = &ticket.reservation_id;
147 let reservation = ctx
148 .store
149 .get(ctx.partition, &keys::reservation(rid)?)
150 .await?
151 .ok_or_else(|| StoreError::Corrupt("ticket has no reservation".into()))?;
152 if !matches!(
153 codec::decode_reservation(&reservation)?,
154 ReservationV1::Ticketed { ticket_id } if ticket_id == id
155 ) {
156 return Err(StoreError::Corrupt("ticket reservation mismatch".into()));
157 }
158 let index = keys::ticket_index(
159 &ticket.repo,
160 &ticket.ref_name,
161 &ticket.pack_id,
162 &ticket.signer,
163 )?;
164 let tc = keys::tickets_per_ref(&ticket.repo, &ticket.ref_name)?;
165 let tu = keys::tickets_per_signer(&ticket.repo, &ticket.ref_name, &ticket.signer)?;
166 let (index_value, ref_count, signer_count, os, oc) = (
167 ctx.store.get(ctx.partition, &index).await?,
168 ctx.store.get(ctx.partition, &tc).await?,
169 ctx.store.get(ctx.partition, &tu).await?,
170 ctx.store
171 .get(ctx.partition, &keys::outbox_sequence())
172 .await?,
173 ctx.store
174 .get(ctx.partition, &keys::outcome_backlog())
175 .await?,
176 );
177 let mut batch = Batch::new();
178 tickets::plan_ticket_close(
179 &id,
180 &ticket,
181 &raw,
182 index_value.as_ref(),
183 ref_count.as_ref(),
184 signer_count.as_ref(),
185 CloseReason::ExpiryTimerFired,
186 &mut batch.preconditions,
187 &mut batch.writes,
188 )?;
189 let mut outbox = OutboxBuilder::new(os.as_ref(), oc.as_ref())?;
190 outbox.outcome(
191 rid,
192 &reservation,
193 Terminal::new(ReservationV1::Expired {
194 repository: repository(ctx.partition, &ticket)?,
195 occurred_at_ms: ctx.now_ms,
196 })?,
197 );
198 outbox.try_finish(&mut batch.preconditions, &mut batch.writes)?;
199
200 verification_cleanup(ctx, &ticket, &mut batch).await?;
201
202 if let Some(session) = &ticket.upload_session
203 && let Err(error) = blobs.abort(BlobKey::pack(ticket.pack_id), session).await
204 {
205 ABORT_FAILURES.fetch_add(1, Ordering::Relaxed);
206 tracing::warn!(error = %error, ticket_id = %mkit_core::hash::to_hex(&id), "ticket expiry session abort failed");
207 }
208 Ok(Fired::Done(batch))
209 };
210 match fire.await {
211 Err(error @ StoreError::Corrupt(_)) => {
212 tracing::warn!(%error, "corrupt ticket expiry row deferred for repair");
213 Ok(Fired::Reschedule {
214 due_at_ms: ctx.now_ms.saturating_add(60_000),
215 value: timer.value.clone(),
216 batch: Batch::new(),
217 })
218 }
219 other => other,
220 }
221 })
222 }
223}
224
225#[cfg(feature = "test-faults")]
227#[derive(Debug)]
228pub(crate) struct BorrowedTicketExpiry<'a, B> {
229 pub(crate) blobs: &'a B,
230}
231
232#[cfg(feature = "test-faults")]
233impl<S: NamespaceStore, B: MultipartBlobStore> TimerHandler<S> for BorrowedTicketExpiry<'_, B> {
234 fn kind(&self) -> TimerKind {
235 kinds::TICKET_EXPIRY
236 }
237
238 fn max_per_tick(&self) -> Option<u32> {
239 Some(8)
240 }
241
242 fn fire<'a>(
243 &'a self,
244 ctx: &'a TimerCtx<'a, S>,
245 timer: &'a DueTimer,
246 ) -> BoxFuture<'a, Result<Fired, StoreError>> {
247 TicketExpiry::fire_with(self.blobs, ctx, timer)
248 }
249}
250
251#[cfg(all(test, feature = "memory"))]
252#[allow(clippy::unwrap_used)] mod tests {
254 use super::*;
255 use crate::repo::RepoName;
256 use crate::store::{
257 BatchOutcome, BlobBody, BlobMeta, BlobStore, ByteRange, UnsupportedPartSink, Value,
258 };
259 use crate::timers::{TickBudget, TimerRegistry, run_due};
260 use crate::{ManualClock, MemoryBlobStore, MemoryKv};
261
262 fn partition() -> Partition {
263 Partition::Namespace(NamespaceKey::deployment_default())
264 }
265
266 fn ticket(n: usize, expiry: u64, session: Option<Vec<u8>>) -> TicketV1 {
267 TicketV1 {
268 authority_generation: None,
269 repo: RepoName::new("repo").unwrap(),
270 ref_name: format!("refs/heads/branch-{n}"),
271 signer: [1; 32],
272 pack_id: [2; 32],
273 bytes: 8 * 1024 * 1024 + 1,
274 part_size: 8 * 1024 * 1024,
275 expires_at_ms: expiry,
276 created_at_ms: 1,
277 reservation_id: format!("s:{n:064x}"),
278 upload_session: session,
279 }
280 }
281
282 async fn plant(store: &MemoryKv, t: &TicketV1, due: u64) -> Hash {
283 let id = tickets::ticket_id(&t.reservation_id);
284 let batch = Batch::new()
285 .put(keys::ticket(&id), codec::encode_ticket(t))
286 .put(
287 keys::reservation(&t.reservation_id).unwrap(),
288 codec::encode_reservation(&ReservationV1::Ticketed { ticket_id: id }),
289 )
290 .put(
291 keys::ticket_index(&t.repo, &t.ref_name, &t.pack_id, &t.signer).unwrap(),
292 codec::encode_ref_id(&id),
293 )
294 .put(
295 keys::tickets_per_ref(&t.repo, &t.ref_name).unwrap(),
296 codec::encode_u64(1),
297 )
298 .put(
299 keys::tickets_per_signer(&t.repo, &t.ref_name, &t.signer).unwrap(),
300 codec::encode_u64(1),
301 )
302 .put(
303 keys::timer(due, kinds::TICKET_EXPIRY.get(), &id),
304 Value::default(),
305 );
306 assert_eq!(
307 store.apply(&partition(), batch).await.unwrap(),
308 BatchOutcome::Committed
309 );
310 id
311 }
312
313 async fn tick(store: &MemoryKv, blobs: MemoryBlobStore, now: u64) -> crate::timers::RunReport {
314 let clock = ManualClock::new(i64::try_from(now).expect("test time fits i64"));
315 run_due(
316 store,
317 &partition(),
318 &TimerRegistry::new().register(TicketExpiry { blobs }),
319 &clock,
320 now,
321 &TickBudget::default(),
322 )
323 .await
324 .unwrap()
325 }
326
327 #[tokio::test]
328 async fn expired_ticket_closes_once_and_aborts_session() {
329 let store = MemoryKv::default();
330 let blobs = MemoryBlobStore::default();
331 let t = ticket(1, 100, None);
332 let id = tickets::ticket_id(&t.reservation_id);
333 let session = blobs
334 .begin_multipart_for_ticket(BlobKey::pack(t.pack_id), t.bytes, t.part_size, id)
335 .await
336 .unwrap();
337 let t = TicketV1 {
338 upload_session: Some(session),
339 ..t
340 };
341 plant(&store, &t, 100).await;
342 assert_eq!(blobs.multipart_session_count(), 1);
343 let report = tick(&store, blobs.clone(), 100).await;
344 assert_eq!(report.fired, 1);
345 assert_eq!(blobs.multipart_session_count(), 0);
346 for key in [
347 keys::ticket(&id),
348 keys::ticket_index(&t.repo, &t.ref_name, &t.pack_id, &t.signer).unwrap(),
349 keys::tickets_per_ref(&t.repo, &t.ref_name).unwrap(),
350 keys::tickets_per_signer(&t.repo, &t.ref_name, &t.signer).unwrap(),
351 keys::timer(100, kinds::TICKET_EXPIRY.get(), &id),
352 ] {
353 assert!(store.get(&partition(), &key).await.unwrap().is_none());
354 }
355 let outcome = store
356 .get(&partition(), &keys::reservation(&t.reservation_id).unwrap())
357 .await
358 .unwrap()
359 .unwrap();
360 assert_eq!(
361 codec::decode_reservation(&outcome).unwrap(),
362 ReservationV1::Expired {
363 repository: "repo".into(),
364 occurred_at_ms: 100,
365 }
366 );
367 assert_eq!(tick(&store, blobs, 100).await.fired, 0);
368 assert_eq!(
369 codec::decode_backlog(
370 &store
371 .get(&partition(), &keys::outcome_backlog())
372 .await
373 .unwrap()
374 .unwrap()
375 )
376 .unwrap()
377 .rows,
378 1
379 );
380 }
381
382 #[tokio::test]
383 async fn expiry_clears_an_unconsumed_packs_verification_and_kicks_its_job() {
384 use crate::indexed::state::{VerificationV1, encode};
385 for member in [false, true] {
386 let store = MemoryKv::default();
387 let t = ticket(4, 100, None);
388 plant(&store, &t, 100).await;
389 let vs = keys::verification(&t.repo, &t.pack_id);
390 let job = keys::verify_job(&t.repo, &t.pack_id);
391 let mut rows = Batch::new()
392 .put(
393 vs.clone(),
394 encode(&VerificationV1::Verified {
395 pack_len: 1,
396 verified_at_ms: 1,
397 publication: None,
398 }),
399 )
400 .put(job.clone(), Value::default());
401 if member {
402 rows = rows.put(keys::membership(&t.repo, &t.pack_id), Value::default());
403 }
404 store.apply(&partition(), rows).await.unwrap();
405 assert_eq!(tick(&store, MemoryBlobStore::default(), 100).await.fired, 1);
406 assert_eq!(
410 store.get(&partition(), &vs).await.unwrap().is_some(),
411 member
412 );
413 let kick = keys::timer(
414 100,
415 kinds::VERIFY.get(),
416 &crate::indexed::checkpoint::timer_reference(&t.repo, &t.pack_id),
417 );
418 assert!(store.get(&partition(), &kick).await.unwrap().is_some());
419 }
420 let store = MemoryKv::default();
422 let t = ticket(5, 100, None);
423 plant(&store, &t, 100).await;
424 assert_eq!(tick(&store, MemoryBlobStore::default(), 100).await.fired, 1);
425 let (start, end) = keys::class_range(keys::TAG_TIMER);
426 let timers = store
427 .scan(&partition(), &start, &end, None, 10)
428 .await
429 .unwrap()
430 .entries;
431 assert!(!timers.iter().any(|(key, _)| matches!(
432 keys::parse(key),
433 Some(keys::ParsedKey::Timer { kind, .. }) if kind == kinds::VERIFY.get()
434 )));
435 }
436
437 #[tokio::test]
438 async fn early_timer_reschedules_and_consumed_ticket_is_done() {
439 let store = MemoryKv::default();
440 let blobs = MemoryBlobStore::default();
441 let t = ticket(2, 200, None);
442 let id = plant(&store, &t, 100).await;
443 assert_eq!(tick(&store, blobs.clone(), 100).await.fired, 1);
444 assert!(
445 store
446 .get(
447 &partition(),
448 &keys::timer(200, kinds::TICKET_EXPIRY.get(), &id)
449 )
450 .await
451 .unwrap()
452 .is_some()
453 );
454 store
455 .apply(&partition(), Batch::new().delete(keys::ticket(&id)))
456 .await
457 .unwrap();
458 assert_eq!(tick(&store, blobs, 200).await.fired, 1);
459 assert!(
460 store
461 .get(&partition(), &keys::reservation(&t.reservation_id).unwrap())
462 .await
463 .unwrap()
464 .is_some_and(|v| matches!(
465 codec::decode_reservation(&v),
466 Ok(ReservationV1::Ticketed { .. })
467 ))
468 );
469 }
470
471 #[tokio::test]
472 async fn consumption_race_rejects_stale_expiry_batch() {
473 let store = MemoryKv::default();
474 let t = ticket(3, 100, None);
475 let id = plant(&store, &t, 100).await;
476 let timer = DueTimer {
477 due_at_ms: 100,
478 kind: kinds::TICKET_EXPIRY,
479 reference: bytes::Bytes::copy_from_slice(&id),
480 value: Value::default(),
481 };
482 let ctx = TimerCtx {
483 store: &store,
484 partition: &partition(),
485 now_ms: 100,
486 };
487 let Fired::Done(stale) = (TicketExpiry {
488 blobs: MemoryBlobStore::default(),
489 })
490 .fire(&ctx, &timer)
491 .await
492 .unwrap() else {
493 panic!("due ticket must close");
494 };
495 let committed = codec::encode_reservation(&ReservationV1::Expired {
496 repository: "repo".into(),
497 occurred_at_ms: 99,
498 });
499 store
500 .apply(
501 &partition(),
502 Batch::new().delete(keys::ticket(&id)).put(
503 keys::reservation(&t.reservation_id).unwrap(),
504 committed.clone(),
505 ),
506 )
507 .await
508 .unwrap();
509 assert!(matches!(
510 store.apply(&partition(), stale).await.unwrap(),
511 BatchOutcome::PreconditionFailed { .. }
512 ));
513 assert_eq!(tick(&store, MemoryBlobStore::default(), 100).await.fired, 1);
514 assert_eq!(
515 store
516 .get(&partition(), &keys::reservation(&t.reservation_id).unwrap())
517 .await
518 .unwrap(),
519 Some(committed)
520 );
521 }
522
523 #[tokio::test]
524 async fn only_eight_tickets_fire_per_tick() {
525 let store = MemoryKv::default();
526 for n in 0..10 {
527 plant(&store, &ticket(n, 100, None), 100).await;
528 }
529 let report = tick(&store, MemoryBlobStore::default(), 100).await;
530 assert_eq!(report.fired, 8);
531 assert_eq!(report.deferred, 2);
532 }
533
534 struct FailingAbort(MemoryBlobStore);
535
536 impl BlobStore for FailingAbort {
537 type Sink = <MemoryBlobStore as BlobStore>::Sink;
538
539 async fn begin(&self, key: BlobKey, len: u64) -> Result<Self::Sink, StoreError> {
540 self.0.begin(key, len).await
541 }
542
543 async fn get(
544 &self,
545 key: &BlobKey,
546 range: Option<ByteRange>,
547 ) -> Result<Option<BlobBody>, StoreError> {
548 self.0.get(key, range).await
549 }
550
551 async fn head(&self, key: &BlobKey) -> Result<Option<BlobMeta>, StoreError> {
552 self.0.head(key).await
553 }
554
555 async fn probe(&self) -> Result<(), StoreError> {
556 self.0.probe().await
557 }
558
559 async fn delete(&self, key: &BlobKey) -> Result<bool, StoreError> {
560 self.0.delete(key).await
561 }
562 }
563
564 impl MultipartBlobStore for FailingAbort {
565 type PartSink = UnsupportedPartSink;
566 const MAX_PARTS: u32 = 1;
567
568 async fn abort(&self, _key: BlobKey, _session: &[u8]) -> Result<(), StoreError> {
569 Err(StoreError::Unavailable("injected abort failure".into()))
570 }
571 }
572
573 #[tokio::test]
574 async fn abort_failure_is_counted_and_does_not_block_expiry() {
575 let store = MemoryKv::default();
576 let t = ticket(11, 100, Some(vec![7; 32]));
577 plant(&store, &t, 100).await;
578 let before = abort_failures();
579 let clock = ManualClock::new(100);
580 let report = run_due(
581 &store,
582 &partition(),
583 &TimerRegistry::new().register(TicketExpiry {
584 blobs: FailingAbort(MemoryBlobStore::default()),
585 }),
586 &clock,
587 100,
588 &TickBudget::default(),
589 )
590 .await
591 .unwrap();
592 assert_eq!(report.fired, 1);
593 assert!(abort_failures() > before);
594 assert!(matches!(
595 codec::decode_reservation(
596 &store
597 .get(&partition(), &keys::reservation(&t.reservation_id).unwrap())
598 .await
599 .unwrap()
600 .unwrap()
601 ),
602 Ok(ReservationV1::Expired { .. })
603 ));
604 }
605
606 #[tokio::test]
607 async fn corrupt_expiry_backs_off_without_starving_healthy_tickets() {
610 let store = MemoryKv::default();
611 let bad = [0_u8; 32];
612 store
613 .apply(
614 &partition(),
615 Batch::new()
616 .put(keys::ticket(&bad), Value::new(&b"bad ticket"[..]))
617 .put(
618 keys::timer(99, kinds::TICKET_EXPIRY.get(), &bad),
619 Value::default(),
620 ),
621 )
622 .await
623 .unwrap();
624 let mut tickets = Vec::new();
625 for n in 20..29 {
626 let t = ticket(n, 100, None);
627 plant(&store, &t, 100).await;
628 tickets.push(t);
629 }
630 let report = tick(&store, MemoryBlobStore::default(), 100).await;
631 assert_eq!(report.failed, 0);
632 assert_eq!(report.fired, 8);
633 assert_eq!(report.deferred, 2);
634 assert_eq!(tick(&store, MemoryBlobStore::default(), 100).await.fired, 2);
635 let closed = futures::future::join_all(tickets.iter().map(|t| async {
636 let value = store
637 .get(&partition(), &keys::reservation(&t.reservation_id).unwrap())
638 .await
639 .unwrap()
640 .unwrap();
641 matches!(
642 codec::decode_reservation(&value),
643 Ok(ReservationV1::Expired { .. })
644 )
645 }))
646 .await;
647 assert_eq!(closed.into_iter().filter(|yes| *yes).count(), 9);
648 }
649}