1use common_telemetry::{debug, warn};
16use puffin::puffin_manager::{PuffinManager, PuffinWriter};
17use store_api::storage::ColumnId;
18
19use crate::sst::file::{RegionFileId, RegionIndexId};
20use crate::sst::index::puffin_manager::SstPuffinWriter;
21use crate::sst::index::statistics::{ByteCount, RowCount};
22use crate::sst::index::{
23 BloomFilterOutput, FulltextIndexOutput, IndexOutput, Indexer, InvertedIndexOutput,
24};
25
26impl Indexer {
27 pub(crate) async fn do_finish(&mut self) -> IndexOutput {
28 self.dense_pk_decoder = None;
29 let mut output = IndexOutput::default();
30
31 let Some(mut writer) = self.build_puffin_writer().await else {
32 self.do_abort().await;
33 return output;
34 };
35
36 let success = self
37 .do_finish_inverted_index(&mut writer, &mut output)
38 .await;
39 if !success {
40 self.do_abort().await;
41 return IndexOutput::default();
42 }
43
44 let success = self
45 .do_finish_fulltext_index(&mut writer, &mut output)
46 .await;
47 if !success {
48 self.do_abort().await;
49 return IndexOutput::default();
50 }
51
52 let success = self.do_finish_bloom_filter(&mut writer, &mut output).await;
53 if !success {
54 self.do_abort().await;
55 return IndexOutput::default();
56 }
57
58 self.do_prune_intm_sst_dir().await;
59 output.file_size = self.do_finish_puffin_writer(writer).await;
60 output.version = self.index_version;
61 output
62 }
63
64 async fn build_puffin_writer(&mut self) -> Option<SstPuffinWriter> {
65 let puffin_manager = self.puffin_manager.clone()?;
66
67 let err = match puffin_manager
68 .writer(&RegionIndexId::new(
69 RegionFileId::new(self.physical_region_id, self.file_id),
70 self.index_version,
71 ))
72 .await
73 {
74 Ok(writer) => return Some(writer),
75 Err(err) => err,
76 };
77
78 if cfg!(any(test, feature = "test")) {
79 panic!(
80 "Failed to create puffin writer, region_id: {}, file_id: {}, err: {:?}",
81 self.region_id, self.file_id, err
82 );
83 } else {
84 warn!(
85 err; "Failed to create puffin writer, region_id: {}, file_id: {}",
86 self.region_id, self.file_id,
87 );
88 }
89
90 None
91 }
92
93 async fn do_finish_puffin_writer(&mut self, writer: SstPuffinWriter) -> ByteCount {
94 let err = match writer.finish().await {
95 Ok(size) => return size,
96 Err(err) => err,
97 };
98
99 if cfg!(any(test, feature = "test")) {
100 panic!(
101 "Failed to finish puffin writer, region_id: {}, file_id: {}, err: {:?}",
102 self.region_id, self.file_id, err
103 );
104 } else {
105 warn!(
106 err; "Failed to finish puffin writer, region_id: {}, file_id: {}",
107 self.region_id, self.file_id,
108 );
109 }
110
111 0
112 }
113
114 async fn do_finish_inverted_index(
116 &mut self,
117 puffin_writer: &mut SstPuffinWriter,
118 index_output: &mut IndexOutput,
119 ) -> bool {
120 let Some(mut indexer) = self.inverted_indexer.take() else {
121 return true;
122 };
123
124 let column_ids = indexer.column_ids().collect();
125 let err = match indexer.finish(puffin_writer).await {
126 Ok((row_count, byte_count)) => {
127 self.fill_inverted_index_output(
128 &mut index_output.inverted_index,
129 row_count,
130 byte_count,
131 column_ids,
132 );
133 return true;
134 }
135 Err(err) => err,
136 };
137
138 if cfg!(any(test, feature = "test")) {
139 panic!(
140 "Failed to finish inverted index, region_id: {}, file_id: {}, err: {:?}",
141 self.region_id, self.file_id, err
142 );
143 } else {
144 warn!(
145 err; "Failed to finish inverted index, region_id: {}, file_id: {}",
146 self.region_id, self.file_id,
147 );
148 }
149
150 false
151 }
152
153 async fn do_finish_fulltext_index(
154 &mut self,
155 puffin_writer: &mut SstPuffinWriter,
156 index_output: &mut IndexOutput,
157 ) -> bool {
158 let Some(mut indexer) = self.fulltext_indexer.take() else {
159 return true;
160 };
161
162 let column_ids = indexer.column_ids().collect();
163 let err = match indexer.finish(puffin_writer).await {
164 Ok((row_count, byte_count)) => {
165 self.fill_fulltext_index_output(
166 &mut index_output.fulltext_index,
167 row_count,
168 byte_count,
169 column_ids,
170 );
171 return true;
172 }
173 Err(err) => err,
174 };
175
176 if cfg!(any(test, feature = "test")) {
177 panic!(
178 "Failed to finish full-text index, region_id: {}, file_id: {}, err: {:?}",
179 self.region_id, self.file_id, err
180 );
181 } else {
182 warn!(
183 err; "Failed to finish full-text index, region_id: {}, file_id: {}",
184 self.region_id, self.file_id,
185 );
186 }
187
188 false
189 }
190
191 async fn do_finish_bloom_filter(
192 &mut self,
193 puffin_writer: &mut SstPuffinWriter,
194 index_output: &mut IndexOutput,
195 ) -> bool {
196 let Some(mut indexer) = self.bloom_filter_indexer.take() else {
197 return true;
198 };
199
200 let column_ids = indexer.column_ids().collect();
201 let err = match indexer.finish(puffin_writer).await {
202 Ok((row_count, byte_count)) => {
203 self.fill_bloom_filter_output(
204 &mut index_output.bloom_filter,
205 row_count,
206 byte_count,
207 column_ids,
208 );
209 return true;
210 }
211 Err(err) => err,
212 };
213
214 if cfg!(any(test, feature = "test")) {
215 panic!(
216 "Failed to finish bloom filter, region_id: {}, file_id: {}, err: {:?}",
217 self.region_id, self.file_id, err
218 );
219 } else {
220 warn!(
221 err; "Failed to finish bloom filter, region_id: {}, file_id: {}",
222 self.region_id, self.file_id,
223 );
224 }
225
226 false
227 }
228
229 fn fill_inverted_index_output(
230 &mut self,
231 output: &mut InvertedIndexOutput,
232 row_count: RowCount,
233 byte_count: ByteCount,
234 column_ids: Vec<ColumnId>,
235 ) {
236 debug!(
237 "Inverted index created, region_id: {}, file_id: {}, written_bytes: {}, written_rows: {}, columns: {:?}",
238 self.region_id, self.file_id, byte_count, row_count, column_ids
239 );
240
241 output.index_size = byte_count;
242 output.row_count = row_count;
243 output.columns = column_ids;
244 }
245
246 fn fill_fulltext_index_output(
247 &mut self,
248 output: &mut FulltextIndexOutput,
249 row_count: RowCount,
250 byte_count: ByteCount,
251 column_ids: Vec<ColumnId>,
252 ) {
253 debug!(
254 "Full-text index created, region_id: {}, file_id: {}, written_bytes: {}, written_rows: {}, columns: {:?}",
255 self.region_id, self.file_id, byte_count, row_count, column_ids
256 );
257
258 output.index_size = byte_count;
259 output.row_count = row_count;
260 output.columns = column_ids;
261 }
262
263 fn fill_bloom_filter_output(
264 &mut self,
265 output: &mut BloomFilterOutput,
266 row_count: RowCount,
267 byte_count: ByteCount,
268 column_ids: Vec<ColumnId>,
269 ) {
270 debug!(
271 "Bloom filter created, region_id: {}, file_id: {}, written_bytes: {}, written_rows: {}, columns: {:?}",
272 self.region_id, self.file_id, byte_count, row_count, column_ids
273 );
274
275 output.index_size = byte_count;
276 output.row_count = row_count;
277 output.columns = column_ids;
278 }
279
280 pub(crate) async fn do_prune_intm_sst_dir(&mut self) {
281 if let Some(manager) = self.intermediate_manager.take()
282 && let Err(e) = manager.prune_sst_dir(&self.region_id, &self.file_id).await
283 {
284 warn!(e; "Failed to prune intermediate SST directory, region_id: {}, file_id: {}", self.region_id, self.file_id);
285 }
286 }
287}