allsource_core/
store_strict_read.rs1use super::{EventStore, ReadScope};
2use crate::{
3 domain::{
4 entities::Event,
5 value_objects::{EntityId, TenantId},
6 },
7 error::{AllSourceError, Result},
8 infrastructure::persistence::archive_budget::ArchiveReadBudget,
9};
10use chrono::SubsecRound;
11use std::{
12 io,
13 sync::{Arc, atomic::AtomicBool},
14};
15
16const MAX_ENTITY_EVENTS: usize = 1_001;
17const MAX_ENCODED_EVENTS: usize = 2 * 1024 * 1024 - 16 * 1024;
19
20impl EventStore {
21 pub fn query_retained_entity(
27 &self,
28 tenant_id: &str,
29 entity_id: &str,
30 limit: usize,
31 scope: &ReadScope,
32 ) -> Result<(Vec<Event>, usize)> {
33 self.query_retained_entity_cancellable(tenant_id, entity_id, limit, scope, None)
34 }
35
36 pub(crate) fn query_retained_entity_cancellable(
37 &self,
38 tenant_id: &str,
39 entity_id: &str,
40 limit: usize,
41 scope: &ReadScope,
42 cancellation: Option<&Arc<AtomicBool>>,
43 ) -> Result<(Vec<Event>, usize)> {
44 TenantId::new(tenant_id.to_string())?;
45 EntityId::new(entity_id.to_string())?;
46 if !(1..=MAX_ENTITY_EVENTS).contains(&limit) || !scope.permits(entity_id) {
47 return Err(AllSourceError::InvalidInput(
48 "Strict retained read target or limit refused".into(),
49 ));
50 }
51 let budget = ArchiveReadBudget::new(self.strict_archive_limits.clone())
52 .with_cancellation(cancellation.cloned());
53 budget.check()?;
54 self.ensure_tenant_loaded_budgeted(tenant_id, true, cancellation.cloned())?;
55 let _resident = self
56 .cache_residency_gate
57 .try_read_for(budget.remaining()?)
58 .ok_or_else(|| {
59 AllSourceError::StorageError("Strict retained read residency lock timed out".into())
60 })?;
61 if !self.tenant_loader.is_complete(tenant_id) {
62 return Err(AllSourceError::StorageError(
63 "Verified archive was evicted before the retained read".into(),
64 ));
65 }
66 let events = self
67 .events
68 .try_read_for(budget.remaining()?)
69 .ok_or_else(|| {
70 AllSourceError::StorageError("Strict retained read event lock timed out".into())
71 })?;
72 let entries = self
73 .index
74 .get_by_entity_bounded(entity_id, MAX_ENTITY_EVENTS)?;
75 let mut selected = Vec::with_capacity(entries.len());
76 let mut encoded = EncodedBudget {
77 remaining: MAX_ENCODED_EVENTS,
78 work: &budget,
79 };
80 for entry in entries {
81 budget.check()?;
82 let event = events
83 .get(entry.offset)
84 .filter(|event| event.id == entry.event_id && event.entity_id_str() == entity_id)
85 .ok_or_else(|| {
86 AllSourceError::StorageError(
87 "Strict retained read encountered an inconsistent index".into(),
88 )
89 })?;
90 if event.tenant_id_str() != tenant_id {
91 continue;
92 }
93 serde_json::to_writer(&mut encoded, event).map_err(|_| {
96 AllSourceError::StorageError(
97 "Strict retained read encoded data budget exceeded".into(),
98 )
99 })?;
100 selected.push(event);
101 }
102 selected.sort_by(|left, right| {
103 left.timestamp
104 .timestamp_micros()
105 .cmp(&right.timestamp.timestamp_micros())
106 .then_with(|| left.version.cmp(&right.version))
107 });
108 let total = selected.len();
109 let result = selected
110 .into_iter()
111 .take(limit)
112 .cloned()
113 .map(|mut event| {
114 event.timestamp = event.timestamp.trunc_subsecs(6);
115 event
116 })
117 .collect();
118 budget.check()?;
119 self.tenant_loader.touch(tenant_id);
120 Ok((result, total))
121 }
122}
123
124struct EncodedBudget<'a> {
125 remaining: usize,
126 work: &'a ArchiveReadBudget,
127}
128
129impl io::Write for EncodedBudget<'_> {
130 fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
131 self.work
132 .check()
133 .map_err(|error| io::Error::other(error.to_string()))?;
134 self.remaining = self
135 .remaining
136 .checked_sub(bytes.len())
137 .ok_or_else(|| io::Error::other("encoded event budget exceeded"))?;
138 Ok(bytes.len())
139 }
140
141 fn flush(&mut self) -> io::Result<()> {
142 Ok(())
143 }
144}
145
146#[cfg(test)]
147mod tests {
148 use super::*;
149 use crate::store::EventStoreConfig;
150 use std::{sync::atomic::Ordering, time::Duration};
151
152 #[tokio::test(flavor = "current_thread")]
153 async fn snapshot_lease_blocks_eviction_and_observes_cancellation_before_return() {
154 for cancel in [false, true] {
155 let directory = tempfile::TempDir::new().unwrap();
156 let store = Arc::new(EventStore::with_config(EventStoreConfig::with_persistence(
157 directory.path(),
158 )));
159 let event = Event::from_strings(
160 "synthetic.updated".into(),
161 "run".into(),
162 "synthetic".into(),
163 serde_json::json!({}),
164 None,
165 )
166 .unwrap();
167 store.ingest_with_expected_version(&event, Some(0)).unwrap();
168 store.flush_storage().unwrap();
169 let cancellation = Arc::new(AtomicBool::new(false));
170 let lock_store = Arc::clone(&store);
173 let (held_tx, held_rx) = tokio::sync::oneshot::channel();
174 let (release_tx, release_rx) = std::sync::mpsc::channel();
175 let holder = std::thread::spawn(move || {
176 let _events = lock_store.events.write();
177 held_tx.send(()).unwrap();
178 release_rx.recv_timeout(Duration::from_secs(3)).unwrap();
180 });
181 held_rx.await.unwrap();
182 let read_store = Arc::clone(&store);
183 let flag = Arc::clone(&cancellation);
184 let reader = tokio::task::spawn_blocking(move || {
185 read_store.query_retained_entity_cancellable(
186 "synthetic",
187 "run",
188 1001,
189 &ReadScope::unrestricted(),
190 Some(&flag),
191 )
192 });
193 tokio::time::timeout(Duration::from_secs(1), async {
195 while store.cache_residency_gate.try_write().is_some() {
196 tokio::task::yield_now().await;
197 }
198 })
199 .await
200 .unwrap();
201 let evict_store = Arc::clone(&store);
202 let (eviction_tx, eviction_rx) = tokio::sync::oneshot::channel();
203 let eviction = tokio::task::spawn_blocking(move || {
204 eviction_tx.send(()).unwrap();
205 evict_store.evict_tenant("synthetic");
206 });
207 eviction_rx.await.unwrap();
208 if cancel {
209 cancellation.store(true, Ordering::Release);
210 }
211 assert!(!eviction.is_finished());
212 release_tx.send(()).unwrap();
213 let result = reader.await.unwrap();
214 eviction.await.unwrap();
215 holder.join().unwrap();
216 if cancel {
217 assert!(result.unwrap_err().to_string().contains("cancelled"));
218 } else {
219 let (snapshot, total) = result.unwrap();
220 assert_eq!(total, 1);
221 assert_eq!(snapshot[0].id, event.id);
222 }
223 assert_eq!(store.total_events(), 0);
224 }
225 }
226}