Skip to main content

mkit_server/timers/
ticket_expiry.rs

1//! Close due upload tickets and queue one terminal `Expired` outcome.
2
3use 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/// Best-effort session abort failures observed by this process.
22#[must_use]
23pub fn abort_failures() -> u64 {
24    ABORT_FAILURES.load(Ordering::Relaxed)
25}
26
27/// Kind-2 expiry handler. The per-kind cap bounds each native and Worker tick.
28#[derive(Debug)]
29pub struct TicketExpiry<B> {
30    /// The blob store backing the tickets on this server.
31    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
54/// An unconsumed pack's verification state goes with its ticket (R-148); a
55/// scheduled job's rows go with it too, by the job's own timer, which this
56/// kicks (WP-4.8). A member pack keeps `vs`: GC removes it with `m`.
57async 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/// Borrow the pipeline's store for the test-only manual timer tick.
226#[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)] // Unwraps assert fixed test fixtures.
253mod 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            // A member pack keeps `vs` (GC removes it with `m`); an unconsumed
407            // one loses it. Either way a scheduled job's timer is kicked to
408            // delete the job's rows.
409            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        // No verification rows: nothing extra is written.
421        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    // The corrupt row takes one slot on its first tick, then backs off 60 s,
608    // so it cannot starve healthy tickets on later ticks.
609    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}