risingwave_storage/hummock/iterator/
forward_concat.rs1use crate::hummock::SstableIterator;
16use crate::hummock::iterator::concat_inner::ConcatIteratorInner;
17
18pub type ConcatIterator = ConcatIteratorInner<SstableIterator>;
20
21#[cfg(test)]
22mod tests {
23 use std::sync::Arc;
24
25 #[cfg(feature = "failpoints")]
26 use foyer::Hint;
27
28 use super::*;
29 #[cfg(feature = "failpoints")]
30 use crate::hummock::CachePolicy;
31 use crate::hummock::iterator::HummockIterator;
32 use crate::hummock::iterator::test_utils::{
33 TEST_KEYS_COUNT, default_builder_opt_for_test, gen_iterator_test_sstable_info,
34 iterator_test_key_of, iterator_test_value_of, mock_sstable_store,
35 };
36 use crate::hummock::sstable::SstableIteratorReadOptions;
37 #[cfg(feature = "failpoints")]
38 use crate::monitor::StoreLocalStatistic;
39
40 #[tokio::test]
41 async fn test_concat_iterator() {
42 let sstable_store = mock_sstable_store().await;
43 let table0 = gen_iterator_test_sstable_info(
44 0,
45 default_builder_opt_for_test(),
46 |x| x,
47 sstable_store.clone(),
48 TEST_KEYS_COUNT,
49 )
50 .await;
51 let table1 = gen_iterator_test_sstable_info(
52 1,
53 default_builder_opt_for_test(),
54 |x| TEST_KEYS_COUNT + x,
55 sstable_store.clone(),
56 TEST_KEYS_COUNT,
57 )
58 .await;
59 let table2 = gen_iterator_test_sstable_info(
60 2,
61 default_builder_opt_for_test(),
62 |x| TEST_KEYS_COUNT * 2 + x,
63 sstable_store.clone(),
64 TEST_KEYS_COUNT,
65 )
66 .await;
67 let mut iter = ConcatIterator::new(
68 vec![table0, table1, table2],
69 sstable_store,
70 Arc::new(SstableIteratorReadOptions::default()),
71 );
72 let mut i = 0;
73 iter.rewind().await.unwrap();
74
75 while iter.is_valid() {
76 let key = iter.key();
77 let val = iter.value();
78 assert_eq!(key, iterator_test_key_of(i).to_ref());
79 assert_eq!(
80 val.into_user_value().unwrap(),
81 iterator_test_value_of(i).as_slice()
82 );
83 i += 1;
84 iter.next().await.unwrap();
85 if i == TEST_KEYS_COUNT * 3 {
86 assert!(!iter.is_valid());
87 break;
88 }
89 }
90
91 iter.rewind().await.unwrap();
92 let key = iter.key();
93 let val = iter.value();
94 assert_eq!(key, iterator_test_key_of(0).to_ref());
95 assert_eq!(
96 val.into_user_value().unwrap(),
97 iterator_test_value_of(0).as_slice()
98 );
99 }
100
101 #[tokio::test]
102 async fn test_concat_seek() {
103 let sstable_store = mock_sstable_store().await;
104 let table0 = gen_iterator_test_sstable_info(
105 0,
106 default_builder_opt_for_test(),
107 |x| x,
108 sstable_store.clone(),
109 TEST_KEYS_COUNT,
110 )
111 .await;
112 let table1 = gen_iterator_test_sstable_info(
113 1,
114 default_builder_opt_for_test(),
115 |x| TEST_KEYS_COUNT + x,
116 sstable_store.clone(),
117 TEST_KEYS_COUNT,
118 )
119 .await;
120 let table2 = gen_iterator_test_sstable_info(
121 2,
122 default_builder_opt_for_test(),
123 |x| TEST_KEYS_COUNT * 2 + x,
124 sstable_store.clone(),
125 TEST_KEYS_COUNT,
126 )
127 .await;
128 let mut iter = ConcatIterator::new(
129 vec![table0, table1, table2],
130 sstable_store,
131 Arc::new(SstableIteratorReadOptions::default()),
132 );
133
134 iter.seek(iterator_test_key_of(TEST_KEYS_COUNT + 1).to_ref())
135 .await
136 .unwrap();
137
138 let key = iter.key();
139 let val = iter.value();
140 assert_eq!(key, iterator_test_key_of(TEST_KEYS_COUNT + 1).to_ref());
141 assert_eq!(
142 val.into_user_value().unwrap(),
143 iterator_test_value_of(TEST_KEYS_COUNT + 1).as_slice()
144 );
145
146 iter.seek(iterator_test_key_of(0).to_ref()).await.unwrap();
148 let key = iter.key();
149 let val = iter.value();
150 assert_eq!(key, iterator_test_key_of(0).to_ref());
151 assert_eq!(
152 val.into_user_value().unwrap(),
153 iterator_test_value_of(0).as_slice()
154 );
155
156 iter.seek(iterator_test_key_of(3 * TEST_KEYS_COUNT - 1).to_ref())
158 .await
159 .unwrap();
160
161 let key = iter.key();
162 let val = iter.value();
163 assert_eq!(key, iterator_test_key_of(3 * TEST_KEYS_COUNT - 1).to_ref());
164 assert_eq!(
165 val.into_user_value().unwrap(),
166 iterator_test_value_of(3 * TEST_KEYS_COUNT - 1).as_slice()
167 );
168
169 iter.seek(iterator_test_key_of(3 * TEST_KEYS_COUNT).to_ref())
171 .await
172 .unwrap();
173 assert!(!iter.is_valid());
174 }
175
176 #[tokio::test]
177 async fn test_concat_seek_not_exists() {
178 let sstable_store = mock_sstable_store().await;
179 let table0 = gen_iterator_test_sstable_info(
180 0,
181 default_builder_opt_for_test(),
182 |x| x * 2,
183 sstable_store.clone(),
184 TEST_KEYS_COUNT,
185 )
186 .await;
187 let table1 = gen_iterator_test_sstable_info(
188 1,
189 default_builder_opt_for_test(),
190 |x| (TEST_KEYS_COUNT + x) * 2,
191 sstable_store.clone(),
192 TEST_KEYS_COUNT,
193 )
194 .await;
195 let table2 = gen_iterator_test_sstable_info(
196 2,
197 default_builder_opt_for_test(),
198 |x| (2 * TEST_KEYS_COUNT + x) * 2,
199 sstable_store.clone(),
200 TEST_KEYS_COUNT,
201 )
202 .await;
203 let mut iter = ConcatIterator::new(
204 vec![table0, table1, table2],
205 sstable_store,
206 Arc::new(SstableIteratorReadOptions::default()),
207 );
208
209 iter.seek(iterator_test_key_of(TEST_KEYS_COUNT + 1).to_ref())
210 .await
211 .unwrap();
212
213 let key = iter.key();
214 let val = iter.value();
215 assert_eq!(key, iterator_test_key_of(TEST_KEYS_COUNT + 2).to_ref());
216 assert_eq!(
217 val.into_user_value().unwrap(),
218 iterator_test_value_of(TEST_KEYS_COUNT + 2).as_slice()
219 );
220
221 iter.seek(iterator_test_key_of((TEST_KEYS_COUNT + 9) * 2 + 1).to_ref())
223 .await
224 .unwrap();
225
226 let key = iter.key();
227 let val = iter.value();
228 assert_eq!(key, iterator_test_key_of(TEST_KEYS_COUNT * 4).to_ref());
229 assert_eq!(
230 val.into_user_value().unwrap(),
231 iterator_test_value_of(TEST_KEYS_COUNT * 4).as_slice()
232 );
233 }
234
235 #[tokio::test]
236 #[cfg(feature = "failpoints")]
237 async fn test_concat_init_error_retains_current_iterator_and_candidate_stats() {
238 let sstable_store = mock_sstable_store().await;
239 let table0 = gen_iterator_test_sstable_info(
240 0,
241 default_builder_opt_for_test(),
242 |x| x,
243 sstable_store.clone(),
244 TEST_KEYS_COUNT,
245 )
246 .await;
247 let table1 = gen_iterator_test_sstable_info(
248 1,
249 default_builder_opt_for_test(),
250 |x| TEST_KEYS_COUNT + x,
251 sstable_store.clone(),
252 TEST_KEYS_COUNT,
253 )
254 .await;
255
256 let mut warmup_stats = StoreLocalStatistic::default();
258 let sstable = sstable_store
259 .sstable(&table1, &mut warmup_stats)
260 .await
261 .unwrap();
262 sstable_store
263 .get(
264 &sstable,
265 0,
266 CachePolicy::Fill(Hint::Normal),
267 &mut warmup_stats,
268 )
269 .await
270 .unwrap();
271 warmup_stats.discard();
272
273 let mut iter = ConcatIterator::new(
274 vec![table0, table1],
275 sstable_store,
276 Arc::new(SstableIteratorReadOptions::default()),
277 );
278 iter.rewind().await.unwrap();
279
280 let mut before = StoreLocalStatistic::default();
281 iter.collect_local_statistic(&mut before);
282
283 fail::cfg("concat_iter_after_seek_before_init_complete", "return").unwrap();
284 let result = iter
285 .seek(iterator_test_key_of(TEST_KEYS_COUNT).to_ref())
286 .await;
287 fail::remove("concat_iter_after_seek_before_init_complete");
288
289 assert!(result.is_err());
290 assert!(iter.is_valid());
291 assert_eq!(iter.key(), iterator_test_key_of(0).to_ref());
292
293 let mut after = StoreLocalStatistic::default();
294 iter.collect_local_statistic(&mut after);
295 assert_eq!(
296 after.cache_data_block_total,
297 before.cache_data_block_total + 1,
298 "failed candidate initialization must settle its cache-hit statistic"
299 );
300 }
301}