Skip to main content

mkit_server/timers/
reservation_reconcile.rs

1//! Reconcile a pending write or read after its safe abandonment deadline.
2
3use crate::rt::BoxFuture;
4use crate::store::codec::{self, AbortReason, ReservationV1};
5use crate::store::outbox::{OutboxBuilder, Terminal};
6use crate::store::{Batch, NamespaceStore, StoreError, keys};
7use crate::timers::{DueTimer, Fired, TimerCtx, TimerHandler, TimerKind, registry::kinds};
8
9/// Kind-9 pending reservation reconciler.
10#[derive(Debug, Default, Clone, Copy)]
11pub struct ReservationReconcile;
12
13impl<S: NamespaceStore> TimerHandler<S> for ReservationReconcile {
14    fn kind(&self) -> TimerKind {
15        kinds::RESERVATION_RECONCILE
16    }
17
18    fn fire<'a>(
19        &'a self,
20        ctx: &'a TimerCtx<'a, S>,
21        timer: &'a DueTimer,
22    ) -> BoxFuture<'a, Result<Fired, StoreError>> {
23        Box::pin(async move {
24            let rid = core::str::from_utf8(&timer.reference)
25                .map_err(|_| StoreError::Corrupt("invalid reconcile reservation id".into()))?;
26            let key = keys::reservation(rid)?;
27            let Some(prior) = ctx.store.get(ctx.partition, &key).await? else {
28                return Ok(Fired::Done(Batch::new()));
29            };
30            let ReservationV1::Pending {
31                repository,
32                reconcile_at_ms,
33                ..
34            } = codec::decode_reservation(&prior)?
35            else {
36                return Ok(Fired::Done(Batch::new()));
37            };
38            if ctx.now_ms < reconcile_at_ms {
39                return Ok(Fired::Reschedule {
40                    due_at_ms: reconcile_at_ms,
41                    value: timer.value.clone(),
42                    batch: Batch::new(),
43                });
44            }
45            let keys = [keys::outbox_sequence(), keys::outcome_backlog()];
46            let values = ctx.store.get_many(ctx.partition, &keys).await?;
47            let mut outbox = OutboxBuilder::new(
48                values.first().and_then(Option::as_ref),
49                values.get(1).and_then(Option::as_ref),
50            )?;
51            outbox.outcome(
52                rid,
53                &prior,
54                Terminal::new(ReservationV1::Aborted {
55                    repository,
56                    occurred_at_ms: ctx.now_ms,
57                    reason: AbortReason::Abandoned,
58                    detail: String::new(),
59                })?,
60            );
61            let mut batch = Batch::new();
62            outbox.try_finish(&mut batch.preconditions, &mut batch.writes)?;
63            Ok(Fired::Done(batch))
64        })
65    }
66}
67
68#[cfg(test)]
69mod tests {
70    use super::*;
71    use crate::memory::MemoryKv;
72    use crate::repo::NamespaceKey;
73    use crate::rt::ManualClock;
74    use crate::store::codec::PendingOp;
75    use crate::store::outbox::OutboxBuilder;
76    use crate::store::{BatchOutcome, Partition, Precondition, Value};
77    use crate::timers::{TickBudget, TimerRegistry, run_due};
78    use std::sync::Arc;
79
80    #[tokio::test]
81    async fn crash_after_pending_reconciles_and_blocks_late_apply() {
82        for (rid, op, reconcile_at) in [
83            ("write", PendingOp::Write, 41_000),
84            ("read", PendingOp::Read, 70_000),
85        ] {
86            let clock = Arc::new(ManualClock::new(0));
87            let store = MemoryKv::with_clock(clock.clone());
88            let partition = Partition::Namespace(NamespaceKey::deployment_default());
89            let pending = ReservationV1::Pending {
90                repository: "repo".into(),
91                created_at_ms: 0,
92                reconcile_at_ms: reconcile_at,
93                op,
94            };
95            let prior = codec::encode_reservation(&pending);
96            let mut builder = OutboxBuilder::new(None, None).unwrap();
97            builder.pending(rid, None, &pending);
98            let mut batch = Batch::new();
99            builder
100                .try_finish(&mut batch.preconditions, &mut batch.writes)
101                .unwrap();
102            assert_eq!(
103                store.apply(&partition, batch).await.unwrap(),
104                BatchOutcome::Committed
105            );
106            let registry = TimerRegistry::new().register(ReservationReconcile);
107            assert_eq!(
108                run_due(
109                    &store,
110                    &partition,
111                    &registry,
112                    clock.as_ref(),
113                    reconcile_at - 1,
114                    &TickBudget::default()
115                )
116                .await
117                .unwrap()
118                .fired,
119                0
120            );
121            clock.set(i64::try_from(reconcile_at).unwrap());
122            assert_eq!(
123                run_due(
124                    &store,
125                    &partition,
126                    &registry,
127                    clock.as_ref(),
128                    reconcile_at,
129                    &TickBudget::default()
130                )
131                .await
132                .unwrap()
133                .fired,
134                1
135            );
136            let key = keys::reservation(rid).unwrap();
137            let value = store.get(&partition, &key).await.unwrap().unwrap();
138            assert!(matches!(
139                codec::decode_reservation(&value).unwrap(),
140                ReservationV1::Aborted {
141                    reason: AbortReason::Abandoned,
142                    ..
143                }
144            ));
145            let late = Batch::new()
146                .require(Precondition::Equals(key, prior))
147                .put(keys::outcome_backlog(), Value::default());
148            assert!(matches!(
149                store.apply(&partition, late).await.unwrap(),
150                BatchOutcome::PreconditionFailed { .. }
151            ));
152            assert_eq!(
153                codec::decode_backlog(
154                    &store
155                        .get(&partition, &keys::outcome_backlog())
156                        .await
157                        .unwrap()
158                        .unwrap()
159                )
160                .unwrap()
161                .rows,
162                1
163            );
164        }
165    }
166}