mkit_server/timers/
reservation_reconcile.rs1use 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#[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 ®istry,
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 ®istry,
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}