1use std::path::Path;
15
16use rudb_common::{LogicalType, Result};
17use rudb_encoding::sequence::grams;
18
19use crate::graph::BUDGET_FLOOR;
20use crate::section::{self, Attachment};
21use crate::{Catalog, Reader, invalid};
22
23pub const TEXT_GRAMS_SHARE: u64 = 50;
30
31pub const SHORTEST: u64 = 16;
37
38#[derive(Debug, Clone, PartialEq, Eq)]
40pub struct Built {
41 pub column: usize,
43 pub rows: usize,
45 pub text_bytes: u64,
47 pub bytes: usize,
49 pub built: bool,
51}
52
53#[must_use]
58pub fn current(reader: &Reader) -> bool {
59 let table = reader.table();
60 text_columns(reader).all(|column| {
61 table.sections().iter().any(|held| {
62 held.kind == *section::TEXT_GRAMS
63 && u64::try_from(column) == Ok(held.id)
64 && held.current(table.generation())
65 })
66 })
67}
68
69pub fn build_text_grams(path: &Path, table: &str) -> Result<Vec<Built>> {
75 build_text_grams_within(path, table, TEXT_GRAMS_SHARE)
76}
77
78pub fn build_text_grams_within(path: &Path, table: &str, share: u64) -> Result<Vec<Built>> {
91 let reader = Catalog::open(path)?.table(table)?;
92 let columns = text_columns(&reader).collect::<Vec<_>>();
93 if columns.is_empty() {
94 return Ok(Vec::new());
95 }
96 let rows = reader.table().rows();
97 let allowance = (reader.layout().columns_total().saturating_mul(share) / 100).max(BUDGET_FLOOR);
98 let mut report = Vec::with_capacity(columns.len());
99 let mut payloads = Vec::with_capacity(columns.len());
100 for &column in &columns {
101 let mut words = Vec::with_capacity(rows * 8);
102 let mut text_bytes = 0_u64;
103 for part in 0..reader.parts() {
104 let chunk = reader.read(part, &[column])?;
105 let values = chunk.column(0)?;
106 for row in 0..chunk.len() {
107 let text = values.bytes_at(row).unwrap_or_default();
108 text_bytes += text.len() as u64;
109 words.extend_from_slice(&grams(text).to_le_bytes());
110 }
111 }
112 if words.len() != rows * 8 {
113 return Err(invalid("a text sketch's row count differs from its table"));
114 }
115 report.push(Built { column, rows, text_bytes, bytes: words.len(), built: false });
116 payloads.push(words);
117 }
118 let mut order = (0..report.len()).collect::<Vec<_>>();
119 order.sort_by_key(|&at| std::cmp::Reverse(report[at].text_bytes));
120 let mut spent = 0_u64;
121 for at in order {
122 let long = report[at].text_bytes >= SHORTEST.saturating_mul(rows as u64);
123 let cost = report[at].bytes as u64;
124 if long && rows > 0 && spent.saturating_add(cost) <= allowance {
125 spent += cost;
126 report[at].built = true;
127 }
128 }
129 drop(reader);
130 let attachments = report
131 .iter()
132 .zip(&payloads)
133 .map(|(one, words)| {
134 Ok(Attachment {
135 kind: *section::TEXT_GRAMS,
136 id: u64::try_from(one.column).map_err(|_| invalid("column index overflow"))?,
137 flags: 0,
138 header_bytes: if one.built {
139 0
140 } else {
141 u32::try_from(one.bytes).unwrap_or(u32::MAX)
142 },
143 bytes: if one.built { words } else { &[] },
144 })
145 })
146 .collect::<Result<Vec<_>>>()?;
147 crate::attach(path, table, &attachments)?;
148 Ok(report)
149}
150
151#[must_use]
156pub fn text_grams(reader: &Reader, column: usize) -> Option<Vec<u64>> {
157 let table = reader.table();
158 let id = u64::try_from(column).ok()?;
159 let held =
160 table.sections().iter().find(|held| held.kind == *section::TEXT_GRAMS && held.id == id)?;
161 if !held.usable(table.generation()) {
162 return None;
163 }
164 let bytes = reader.payload(held).ok()?;
165 if bytes.len() != table.rows().checked_mul(8)? {
166 return None;
167 }
168 Some(
169 bytes
170 .chunks_exact(8)
171 .map(|word| u64::from_le_bytes(word.try_into().unwrap_or_default()))
172 .collect(),
173 )
174}
175
176fn text_columns(reader: &Reader) -> impl Iterator<Item = usize> + '_ {
177 reader
178 .table()
179 .fields()
180 .iter()
181 .enumerate()
182 .filter(|(_, field)| field.ty == LogicalType::Varchar)
183 .map(|(column, _)| column)
184}
185
186#[cfg(test)]
187mod tests {
188 use std::fs;
189 use std::path::PathBuf;
190 use std::time::{SystemTime, UNIX_EPOCH};
191
192 use rudb_common::{Field, Value};
193 use rudb_encoding::sequence::Sequence;
194 use rudb_vector::{Chunk, Vector};
195
196 use super::*;
197 use crate::Writer;
198
199 fn path(label: &str) -> PathBuf {
200 let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
201 std::env::temp_dir().join(format!("rudb-grams-{label}-{}-{stamp}.rdb", std::process::id()))
202 }
203
204 fn table_of(label: &str, rows: usize) -> (PathBuf, Vec<String>) {
206 let path = path(label);
207 let words = [
208 "special",
209 "requests",
210 "pending",
211 "deposits",
212 "carefully",
213 "final",
214 "ironic",
215 "slyly",
216 "quickly",
217 "packages",
218 "accounts",
219 "furiously",
220 "express",
221 "regular",
222 "bold",
223 "even",
224 ];
225 let mut seed = 11_u64;
226 let comments = (0..rows)
227 .map(|_| {
228 seed = seed
229 .wrapping_mul(6_364_136_223_846_793_005)
230 .wrapping_add(1_442_695_040_888_963_407);
231 (0..8)
232 .map(|at| words[(seed >> (8 + at * 6)) as usize % words.len()])
233 .collect::<Vec<_>>()
234 .join(" ")
235 })
236 .collect::<Vec<_>>();
237 let mut writer = Writer::create(
238 &path,
239 "orders",
240 vec![
241 Field::new("comment", LogicalType::Varchar),
242 Field::new("flag", LogicalType::Varchar),
243 ],
244 )
245 .expect("new file");
246 for (at, part) in comments.chunks(1000).enumerate() {
247 let text = part.iter().map(|one| Value::Varchar(one.clone())).collect::<Vec<_>>();
248 let flags = part
249 .iter()
250 .map(|_| Value::Varchar(if at % 2 == 0 { "F" } else { "O" }.into()))
251 .collect::<Vec<_>>();
252 let chunk = Chunk::new(vec![
253 Vector::from_values(LogicalType::Varchar, &text).expect("comments"),
254 Vector::from_values(LogicalType::Varchar, &flags).expect("flags"),
255 ])
256 .expect("two columns");
257 writer.append(&chunk).expect("a part");
258 }
259 writer.finish().expect("commit");
260 (path, comments)
261 }
262
263 #[test]
264 fn a_long_column_is_sketched_and_a_short_one_is_recorded_without_bytes() {
265 let (path, comments) = table_of("kept", 3000);
266 let built = build_text_grams(&path, "orders").expect("build");
267 assert_eq!(built.len(), 2);
268 assert!(built[0].built, "a comment column is long enough to sketch");
269 assert!(!built[1].built, "a one byte flag is not");
270
271 let reader = Catalog::open(&path).expect("reopen").table("orders").expect("the table");
272 assert!(current(&reader), "both columns have an entry");
273 assert!(text_grams(&reader, 1).is_none());
274 let sketch = text_grams(&reader, 0).expect("the comment sketch is in the file");
275 assert_eq!(sketch.len(), 3000);
276 for (word, comment) in sketch.iter().zip(&comments) {
277 assert_eq!(*word, grams(comment.as_bytes()));
278 }
279 fs::remove_file(&path).expect("clean up");
280 }
281
282 #[test]
283 fn a_like_answered_through_the_sketch_keeps_the_rows_it_kept_without_one() {
284 let (path, comments) = table_of("like", 6000);
285 let before = Catalog::open(&path).expect("open").table("orders").expect("the table");
286 let sequence = Sequence::new(&[b"special", b"requests"]).expect("an automaton");
287 let mut walked = Vec::new();
288 for part in 0..before.parts() {
289 walked.push(before.rows_holding(part, 0, &sequence, true).expect("a part"));
290 }
291 drop(before);
292 build_text_grams(&path, "orders").expect("build");
293 let after = Catalog::open(&path).expect("reopen").table("orders").expect("the table");
294 let mut kept = 0;
295 for (part, walked) in walked.iter().enumerate() {
296 let sketched = after.rows_holding(part, 0, &sequence, true).expect("a part");
297 assert_eq!(&sketched, walked, "part {part}");
298 kept += sketched.map_or(0, |rows| rows.len());
299 }
300 let wanted = comments
301 .iter()
302 .filter(|one| !one.find("special").is_some_and(|at| one[at + 7..].contains("requests")))
303 .count();
304 assert!(walked.iter().all(Option::is_some), "every part is compressed text");
305 assert!(wanted > 0 && wanted < comments.len());
306 assert_eq!(kept, wanted);
307 fs::remove_file(&path).expect("clean up");
308 }
309}