1use crate::structured_db::{
2 StructuredColumnFamilySchema, StructuredColumnValue, StructuredScanOptions, combined_resolver,
3 decode_row, load_structured_schema_from_cobble_schema,
4};
5use bytes::Bytes;
6use cobble::{
7 Config, GlobalSnapshotManifest, MergeOperatorResolver, ReadOnlyDb, Result, ScanPlan, ScanSplit,
8 ScanSplitScanner, ShardSnapshotRef,
9};
10use serde::{Deserialize, Serialize};
11use std::sync::Arc;
12
13pub struct StructuredScanPlan {
17 inner: ScanPlan,
18}
19
20impl StructuredScanPlan {
21 pub fn new(manifest: GlobalSnapshotManifest) -> Self {
22 Self {
23 inner: ScanPlan::new(manifest),
24 }
25 }
26
27 pub fn with_start(mut self, start: Vec<u8>) -> Self {
28 self.inner = self.inner.with_start(start);
29 self
30 }
31
32 pub fn with_end(mut self, end: Vec<u8>) -> Self {
33 self.inner = self.inner.with_end(end);
34 self
35 }
36
37 pub fn splits(&self) -> Vec<StructuredScanSplit> {
38 self.inner.splits().into_iter().map(Into::into).collect()
39 }
40}
41
42#[derive(Clone, Debug, Serialize, Deserialize)]
44pub struct StructuredScanSplit {
45 pub shard: ShardSnapshotRef,
46 pub start: Option<Vec<u8>>,
47 pub end: Option<Vec<u8>>,
48 pub start_bucket: Option<u16>,
49 pub start_key_exclusive: Option<Vec<u8>>,
50 pub end_bucket: Option<u16>,
51 pub end_key_inclusive: Option<Vec<u8>>,
52}
53
54pub struct StructuredScanSplitPartition {
55 pub before: StructuredScanSplit,
56 pub after: StructuredScanSplit,
57}
58
59impl From<ScanSplit> for StructuredScanSplit {
60 fn from(value: ScanSplit) -> Self {
61 Self {
62 shard: value.shard,
63 start: value.start,
64 end: value.end,
65 start_bucket: value.start_bucket,
66 start_key_exclusive: value.start_key_exclusive,
67 end_bucket: value.end_bucket,
68 end_key_inclusive: value.end_key_inclusive,
69 }
70 }
71}
72
73impl From<StructuredScanSplit> for ScanSplit {
74 fn from(value: StructuredScanSplit) -> Self {
75 Self {
76 shard: value.shard,
77 start: value.start,
78 end: value.end,
79 start_bucket: value.start_bucket,
80 start_key_exclusive: value.start_key_exclusive,
81 end_bucket: value.end_bucket,
82 end_key_inclusive: value.end_key_inclusive,
83 }
84 }
85}
86
87impl StructuredScanSplit {
88 pub fn split_after(
89 &self,
90 bucket: u16,
91 key_inclusive: Vec<u8>,
92 ) -> Result<StructuredScanSplitPartition> {
93 let partition = ScanSplit::from(self.clone()).split_after(bucket, key_inclusive)?;
94 Ok(StructuredScanSplitPartition {
95 before: partition.before.into(),
96 after: partition.after.into(),
97 })
98 }
99
100 pub fn create_scanner_without_options(
101 &self,
102 config: Config,
103 ) -> Result<StructuredScanSplitScanner> {
104 self.create_scanner_without_options_internal(config, None)
105 }
106
107 pub fn create_scanner(
108 &self,
109 config: Config,
110 options: &StructuredScanOptions,
111 ) -> Result<StructuredScanSplitScanner> {
112 self.create_scanner_internal(config, None, options)
113 }
114
115 pub fn create_scanner_with_resolver_without_options(
116 &self,
117 config: Config,
118 resolver: Arc<dyn MergeOperatorResolver>,
119 ) -> Result<StructuredScanSplitScanner> {
120 self.create_scanner_without_options_internal(config, Some(resolver))
121 }
122
123 pub fn create_scanner_with_resolver(
124 &self,
125 config: Config,
126 resolver: Arc<dyn MergeOperatorResolver>,
127 options: &StructuredScanOptions,
128 ) -> Result<StructuredScanSplitScanner> {
129 self.create_scanner_internal(config, Some(resolver), options)
130 }
131
132 fn create_scanner_internal(
133 &self,
134 config: Config,
135 resolver: Option<Arc<dyn MergeOperatorResolver>>,
136 options: &StructuredScanOptions,
137 ) -> Result<StructuredScanSplitScanner> {
138 let resolver = combined_resolver(resolver);
139 let read_only = ReadOnlyDb::open_with_db_id_and_resolver(
140 config.clone(),
141 self.shard.snapshot_id,
142 self.shard.db_id.clone(),
143 Arc::clone(&resolver),
144 )?;
145 let structured_schema = Arc::new(load_structured_schema_from_cobble_schema(
146 &read_only.current_schema(),
147 )?);
148 let projected_schema = options.resolve_projected_schema_cached(&structured_schema)?;
149 let scanner = ScanSplit::from(self.clone()).create_scanner_with_resolver(
150 config,
151 resolver,
152 options.as_cobble(),
153 )?;
154 Ok(StructuredScanSplitScanner {
155 inner: scanner,
156 structured_schema: projected_schema,
157 })
158 }
159
160 fn create_scanner_without_options_internal(
161 &self,
162 config: Config,
163 resolver: Option<Arc<dyn MergeOperatorResolver>>,
164 ) -> Result<StructuredScanSplitScanner> {
165 let resolver = combined_resolver(resolver);
166 let read_only = ReadOnlyDb::open_with_db_id_and_resolver(
167 config.clone(),
168 self.shard.snapshot_id,
169 self.shard.db_id.clone(),
170 Arc::clone(&resolver),
171 )?;
172 let structured_schema = Arc::new(load_structured_schema_from_cobble_schema(
173 &read_only.current_schema(),
174 )?);
175 let projected_schema = Arc::new(structured_schema.projected(0, None));
176 let scanner = ScanSplit::from(self.clone())
177 .create_scanner_with_resolver_without_options(config, resolver)?;
178 Ok(StructuredScanSplitScanner {
179 inner: scanner,
180 structured_schema: projected_schema,
181 })
182 }
183}
184
185pub struct StructuredScanSplitScanner {
186 inner: ScanSplitScanner,
187 structured_schema: Arc<StructuredColumnFamilySchema>,
188}
189
190impl StructuredScanSplitScanner {
191 pub fn consume_next_row<T, F>(&mut self, mut consumer: F) -> Result<Option<T>>
192 where
193 F: FnMut(&Bytes, &[Option<StructuredColumnValue>]) -> Result<T>,
194 {
195 let structured_schema = Arc::clone(&self.structured_schema);
196 self.inner.consume_next_row(|key, columns| {
197 let decoded = decode_row(&structured_schema, 0, columns.to_vec())?;
198 consumer(key, &decoded)
199 })
200 }
201
202 pub fn consume_next_row_with_bucket<T, F>(&mut self, mut consumer: F) -> Result<Option<T>>
203 where
204 F: FnMut(u16, &Bytes, &[Option<StructuredColumnValue>]) -> Result<T>,
205 {
206 let structured_schema = Arc::clone(&self.structured_schema);
207 self.inner
208 .consume_next_row_with_bucket(|bucket, key, columns| {
209 let decoded = decode_row(&structured_schema, 0, columns.to_vec())?;
210 consumer(bucket, key, &decoded)
211 })
212 }
213}
214
215impl Iterator for StructuredScanSplitScanner {
216 type Item = Result<(u16, Bytes, Vec<Option<StructuredColumnValue>>)>;
217
218 fn next(&mut self) -> Option<Self::Item> {
219 self.inner.next().map(|item| {
220 let (bucket, key, columns) = item?;
221 let decoded = decode_row(&self.structured_schema, 0, columns)?;
222 Ok((bucket, key, decoded))
223 })
224 }
225}
226
227#[cfg(test)]
228#[path = "../tests/unit/structured_scan.rs"]
229mod tests;