1use super::{
3 Batch, BatchOutcome, Cursor, Key, NamespaceStore, Partition, PartitionStats, RangeScan,
4 ScanPage, StoreCapabilities, StoreError, Value, keys, publication::Witness,
5};
6use crate::pipeline::clearance::PublicationPolicy;
7use crate::repo::RepoId;
8
9pub struct ViewStore<'a, S> {
11 pub store: &'a S,
13 pub repo: &'a RepoId,
15 pub writer: bool,
17 pub policy: Option<&'a dyn PublicationPolicy>,
19}
20impl<S> core::fmt::Debug for ViewStore<'_, S> {
21 fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
22 f.debug_struct("ViewStore")
23 .field("repo", self.repo)
24 .field("writer", &self.writer)
25 .field("serving_stop", &self.policy.is_some())
26 .finish_non_exhaustive()
27 }
28}
29fn retag(key: &Key, from: &[u8], to: &[u8]) -> Key {
30 let bytes = key.as_bytes();
31 if bytes.starts_with(from) {
32 Key::new([to, &bytes[from.len()..]].concat())
33 } else {
34 key.clone()
35 }
36}
37impl<S: NamespaceStore> ViewStore<'_, S> {
38 fn key(&self, p: &Partition, key: &Key) -> Key {
39 if self.writer || !self.store.capabilities().atomic_multi_key {
40 return key.clone();
41 }
42 let key = retag(key, b"r\0", b"pr\0");
43 let key = retag(&key, b"x\0", b"py\0");
44 if matches!(p, Partition::RepoIndex { .. }) {
45 retag(&key, b"m\0", b"pm\0")
46 } else {
47 key
48 }
49 }
50 fn cursor(&self, cursor: Option<&Cursor>, published: bool) -> Option<Cursor> {
51 cursor.map(|cursor| {
52 let live = retag(
53 &retag(&Key::new(cursor.as_bytes().to_vec()), b"pr\0", b"r\0"),
54 b"py\0",
55 b"x\0",
56 );
57 Cursor::new(
58 if published {
59 self.key(&Partition::Namespace(self.repo.namespace.clone()), &live)
60 } else {
61 live
62 }
63 .into_bytes(),
64 )
65 })
66 }
67 fn value(&self, key: &Key, value: Option<Value>) -> Result<Option<Value>, StoreError> {
68 let Some(value) = value else {
69 return Ok(None);
70 };
71 if let Some(keys::ParsedKey::Membership { repo, pack_id }) = keys::parse(key) {
72 if repo != self.repo.name {
73 return Err(StoreError::Corrupt("view repository mismatch".into()));
74 }
75 let witness = Witness::decode(&value)?;
76 if !witness.visible(self.writer, 0)
77 || self
78 .policy
79 .is_some_and(|p| !p.pack_available(self.repo, &pack_id))
80 {
81 return Ok(None);
82 }
83 }
84 Ok(Some(value))
85 }
86 fn page(&self, mut page: ScanPage) -> ScanPage {
87 if !self.writer && self.store.capabilities().atomic_multi_key {
88 for (key, _) in &mut page.entries {
89 *key = retag(&retag(key, b"pr\0", b"r\0"), b"py\0", b"x\0");
90 }
91 }
92 page.next = self.cursor(page.next.as_ref(), false);
93 page
94 }
95}
96impl<S: NamespaceStore> NamespaceStore for ViewStore<'_, S> {
97 fn capabilities(&self) -> StoreCapabilities {
98 self.store.capabilities()
99 }
100 async fn get(&self, p: &Partition, key: &Key) -> Result<Option<Value>, StoreError> {
101 let mut values = self.get_many(p, core::slice::from_ref(key)).await?;
102 Ok(values.pop().flatten())
103 }
104 async fn get_many(
105 &self,
106 p: &Partition,
107 keys: &[Key],
108 ) -> Result<Vec<Option<Value>>, StoreError> {
109 let selected = keys.iter().map(|key| self.key(p, key)).collect::<Vec<_>>();
110 let values = self.store.get_many(p, &selected).await?;
111 if values.len() != keys.len() {
112 return Err(StoreError::Corrupt("short view read".into()));
113 }
114 keys.iter()
115 .zip(values)
116 .map(|(key, value)| self.value(key, value))
117 .collect()
118 }
119 async fn scan(
120 &self,
121 p: &Partition,
122 start: &Key,
123 end: &Key,
124 after: Option<&Cursor>,
125 limit: u32,
126 ) -> Result<ScanPage, StoreError> {
127 let ranges = [RangeScan {
128 start: start.clone(),
129 end: end.clone(),
130 after: after.cloned(),
131 limit,
132 }];
133 let mut pages = self.scan_many(p, &ranges).await?;
134 pages
135 .pop()
136 .ok_or_else(|| StoreError::Corrupt("short view scan".into()))
137 }
138 async fn scan_many(
139 &self,
140 p: &Partition,
141 ranges: &[RangeScan],
142 ) -> Result<Vec<ScanPage>, StoreError> {
143 if ranges.is_empty() || self.writer || !self.store.capabilities().atomic_multi_key {
144 return self.store.scan_many(p, ranges).await;
145 }
146 let selected = ranges
147 .iter()
148 .map(|range| RangeScan {
149 start: self.key(p, &range.start),
150 end: self.key(p, &range.end),
151 after: self.cursor(range.after.as_ref(), true),
152 limit: range.limit,
153 })
154 .collect::<Vec<_>>();
155 let mut captured = Vec::with_capacity(selected.len());
156 while captured.len() < selected.len() {
157 let pages = self.store.scan_many(p, &selected[captured.len()..]).await?;
158 if pages.is_empty() || pages.len() > selected.len() - captured.len() {
159 return Err(StoreError::Corrupt("short view scan".into()));
160 }
161 captured.extend(pages);
162 }
163 Ok(captured.into_iter().map(|page| self.page(page)).collect())
164 }
165 async fn apply(&self, _: &Partition, _: Batch) -> Result<BatchOutcome, StoreError> {
166 Err(StoreError::Invalid(
167 "read view cannot mutate storage".into(),
168 ))
169 }
170 async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
171 self.store.stats(p).await
172 }
173 async fn probe(&self) -> Result<(), StoreError> {
174 self.store.probe().await
175 }
176}