Skip to main content

allsource_core/
store_strict_read.rs

1use 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;
17// Leave room inside the consumer's 2 MiB wire cap for the response envelope.
18const MAX_ENCODED_EVENTS: usize = 2 * 1024 * 1024 - 16 * 1024;
19
20impl EventStore {
21    /// Read one complete retained entity snapshot under an explicit read scope.
22    /// Archive errors, oversized input and an evicted certification refuse the
23    /// read; this is not evidence that retention never removed older history.
24    /// Timestamp precision and ordering match Parquet's microseconds so the
25    /// same retained evidence stays stable across cache eviction and restart.
26    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            // Count serialized bytes without allocating another payload-sized
94            // buffer. Refuse before cloning any selected events.
95            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            // Hold only the materialization lock: the reader can first acquire
171            // its residency lease, then wait here with its original deadline.
172            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                // test-hang-allow: bounded synthetic contention, including test failures.
179                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            // test-hang-allow: bounded observation of the real reader pinning cache residency.
194            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}