search.rs
| 1 | //! `GET /api/search` — index-free name and content search, streamed as SSE. |
| 2 | //! |
| 3 | //! No index, by design: every search walks the selected root with |
| 4 | //! [`ignore::WalkParallel`] (fd/ripgrep's walker: parallel, skips hidden |
| 5 | //! files, honors `.gitignore`) and matches on the fly. |
| 6 | //! |
| 7 | //! * **Name** (scope `name`/`both`): every word of the query must occur |
| 8 | //! (case-insensitive) in the entry's own name — the last path component, |
| 9 | //! not the whole relative path — so "invoice 2025" matches |
| 10 | //! `invoice_2025-11_final.pdf` but a file is not a hit merely for sitting |
| 11 | //! inside a directory whose name matches. |
| 12 | //! * **Content** (scope `content`/`both`): the whole query as a literal, |
| 13 | //! case-insensitive, ripgrep-style scan via the `grep` crates (ripgrep's |
| 14 | //! engine). Binary files are skipped by NUL detection, files over |
| 15 | //! [`SEARCH_MAX_FILE_BYTES`] are skipped and counted. |
| 16 | //! |
| 17 | //! Both scopes run in one walk, so `both` reads the tree once and its file |
| 18 | //! and match events interleave. |
| 19 | //! |
| 20 | //! Results are one JSON object per SSE event, in the order found; the stream |
| 21 | //! always ends with a `done` event. The search stops when the client goes |
| 22 | //! away: axum drops the response body stream, the channel receiver is |
| 23 | //! dropped with it, and the walkers bail at the next `is_closed` check — |
| 24 | //! no per-search budgets or timeouts. |
| 25 | |
| 26 | use std::path::{Path, PathBuf}; |
| 27 | use std::sync::Arc; |
| 28 | use std::sync::atomic::{AtomicUsize, Ordering}; |
| 29 | use std::time::Instant; |
| 30 | |
| 31 | use api_types::SearchEvent; |
| 32 | use axum::extract::{Query as AxumQuery, State}; |
| 33 | use axum::http::StatusCode; |
| 34 | use axum::response::sse::{Event, Sse}; |
| 35 | use axum::response::{IntoResponse, Response}; |
| 36 | use futures_util::StreamExt; |
| 37 | use grep::regex::RegexMatcherBuilder; |
| 38 | use grep::searcher::{BinaryDetection, Searcher, SearcherBuilder, Sink, SinkMatch}; |
| 39 | use ignore::{DirEntry, WalkBuilder, WalkState}; |
| 40 | use serde::Deserialize; |
| 41 | |
| 42 | use crate::api::common::AuthUser; |
| 43 | use crate::db::RootRow; |
| 44 | use crate::error::{ApiError, AppState}; |
| 45 | |
| 46 | /// Content search skips files larger than this. There is no search budget |
| 47 | /// (the client can stop at any time), but a multi-gigabyte text file would |
| 48 | /// stall its worker indefinitely — ripgrep has no such cap, we do. |
| 49 | const SEARCH_MAX_FILE_BYTES: u64 = 10 * 1024 * 1024; |
| 50 | /// A matched line is truncated to this many bytes on the wire. |
| 51 | const SEARCH_MAX_LINE_BYTES: usize = 500; |
| 52 | /// Total matched lines emitted per search. A one-character query can match |
| 53 | /// tens of thousands of lines; unbounded streaming would flood the client. |
| 54 | const SEARCH_MAX_MATCH_EVENTS: usize = 50_000; |
| 55 | /// Matched lines emitted per file (mirrors the client's render cap). |
| 56 | const SEARCH_MAX_LINES_PER_FILE: usize = 500; |
| 57 | |
| 58 | #[derive(Debug, Deserialize)] |
| 59 | pub(super) struct SearchQuery { |
| 60 | q: Option<String>, |
| 61 | /// `name` (default), `content` or `both`. |
| 62 | #[serde(default)] |
| 63 | scope: Option<String>, |
| 64 | /// The root to search; omitted = the caller's first root. |
| 65 | #[serde(default)] |
| 66 | root: Option<i64>, |
| 67 | } |
| 68 | |
| 69 | pub(super) async fn search( |
| 70 | State(state): State<Arc<AppState>>, |
| 71 | auth: AuthUser, |
| 72 | AxumQuery(query): AxumQuery<SearchQuery>, |
| 73 | ) -> Result<Response, ApiError> { |
| 74 | // Search is a session-user feature: a share visitor has no search UI and |
| 75 | // must not probe files outside the shared item. |
| 76 | if auth.share.is_some() { |
| 77 | return Err(ApiError::localized( |
| 78 | StatusCode::FORBIDDEN, |
| 79 | "search requires a signed-in session", |
| 80 | "err_search_forbidden", |
| 81 | )); |
| 82 | } |
| 83 | let q = query |
| 84 | .q |
| 85 | .as_deref() |
| 86 | .map(str::trim) |
| 87 | .filter(|q| !q.is_empty()) |
| 88 | .ok_or_else(|| { |
| 89 | ApiError::localized( |
| 90 | StatusCode::BAD_REQUEST, |
| 91 | "search requires a non-empty query", |
| 92 | "err_search_empty_query", |
| 93 | ) |
| 94 | })?; |
| 95 | |
| 96 | let scope = query.scope.as_deref().unwrap_or("name"); |
| 97 | let want_name = matches!(scope, "name" | "both"); |
| 98 | let want_content = matches!(scope, "content" | "both"); |
| 99 | if !(want_name || want_content) { |
| 100 | return Err(ApiError::localized( |
| 101 | StatusCode::BAD_REQUEST, |
| 102 | "scope must be name, content or both", |
| 103 | "err_search_bad_scope", |
| 104 | )); |
| 105 | } |
| 106 | |
| 107 | // Only a root the caller actually has. |
| 108 | let root = match query.root { |
| 109 | Some(id) => auth.roots.iter().find(|r| r.id == id), |
| 110 | None => auth.roots.first(), |
| 111 | } |
| 112 | .cloned() |
| 113 | .ok_or_else(|| { |
| 114 | ApiError::localized( |
| 115 | StatusCode::FORBIDDEN, |
| 116 | "root not accessible", |
| 117 | "err_root_forbidden", |
| 118 | ) |
| 119 | })?; |
| 120 | |
| 121 | let events = search_stream( |
| 122 | q.to_string(), |
| 123 | want_name, |
| 124 | want_content, |
| 125 | root, |
| 126 | state.root.clone(), |
| 127 | ); |
| 128 | // `Sse` does the `data: <json>\n\n` framing and the content-type and |
| 129 | // cache headers; nginx and friends still need telling not to buffer. |
| 130 | let sse = Sse::new(events.map(|ev| Event::default().json_data(&ev))); |
| 131 | Ok(([("x-accel-buffering", "no")], sse).into_response()) |
| 132 | } |
| 133 | |
| 134 | /// Shared state of the search, behind one `Arc` that every walk visitor |
| 135 | /// clones a handle to. The counters are atomics so they can be updated from |
| 136 | /// the walker threads lock-free. |
| 137 | struct SearchState { |
| 138 | tx: tokio::sync::mpsc::Sender<SearchEvent>, |
| 139 | words: Vec<String>, |
| 140 | root: RootRow, |
| 141 | server_root: PathBuf, |
| 142 | started: Instant, |
| 143 | /// Directory entries visited. |
| 144 | scanned: AtomicUsize, |
| 145 | /// Files skipped for content search (over the size cap). |
| 146 | skipped: AtomicUsize, |
| 147 | names: AtomicUsize, |
| 148 | matches: AtomicUsize, |
| 149 | } |
| 150 | |
| 151 | /// The search itself: one walk over the root, matching names and grepping |
| 152 | /// contents as it goes, sending events as they are found. Runs on a plain |
| 153 | /// thread — `WalkParallel` manages its own worker pool — while the SSE body |
| 154 | /// stream polls the receiver. |
| 155 | /// |
| 156 | /// Stopping: when the client goes away, axum drops the body stream, which |
| 157 | /// drops the receiver. Every `is_closed` check in the walkers then yields |
| 158 | /// `WalkState::Quit` at the next entry, so the walk unwinds promptly. |
| 159 | fn search_stream( |
| 160 | q: String, |
| 161 | want_name: bool, |
| 162 | want_content: bool, |
| 163 | root: RootRow, |
| 164 | server_root: PathBuf, |
| 165 | ) -> impl futures_util::Stream<Item = SearchEvent> + Send { |
| 166 | let (tx, rx) = tokio::sync::mpsc::channel::<SearchEvent>(256); |
| 167 | let state = Arc::new(SearchState { |
| 168 | tx, |
| 169 | words: q.split_whitespace().map(|w| w.to_lowercase()).collect(), |
| 170 | root, |
| 171 | server_root, |
| 172 | started: Instant::now(), |
| 173 | scanned: AtomicUsize::new(0), |
| 174 | skipped: AtomicUsize::new(0), |
| 175 | names: AtomicUsize::new(0), |
| 176 | matches: AtomicUsize::new(0), |
| 177 | }); |
| 178 | |
| 179 | std::thread::spawn(move || { |
| 180 | walk(&state, &q, want_name, want_content); |
| 181 | // `stopped`: the receiver went away (client stopped or navigated) |
| 182 | // or the match cap was hit, before the walk finished. |
| 183 | let stopped = state.tx.is_closed() |
| 184 | || state.matches.load(Ordering::Relaxed) >= SEARCH_MAX_MATCH_EVENTS; |
| 185 | let _ = state.tx.blocking_send(SearchEvent::Done { |
| 186 | stopped, |
| 187 | files: state.names.load(Ordering::Relaxed), |
| 188 | matches: state.matches.load(Ordering::Relaxed), |
| 189 | scanned: state.scanned.load(Ordering::Relaxed), |
| 190 | skipped: state.skipped.load(Ordering::Relaxed), |
| 191 | elapsed_ms: state.started.elapsed().as_millis() as u64, |
| 192 | }); |
| 193 | }); |
| 194 | |
| 195 | tokio_stream::wrappers::ReceiverStream::new(rx) |
| 196 | } |
| 197 | |
| 198 | /// True when every query word occurs in `name`. Both are already lowercased. |
| 199 | /// |
| 200 | /// `name` is the entry's own name, never its path — see the module docs. |
| 201 | fn name_matches(words: &[String], name: &str) -> bool { |
| 202 | words.iter().all(|w| name.contains(w.as_str())) |
| 203 | } |
| 204 | |
| 205 | /// Walks the root in parallel, matching each entry against the wanted |
| 206 | /// scopes. The visitor is cloned once per worker thread, and returns |
| 207 | /// [`WalkState::Quit`] to stop the whole walk (client went away). |
| 208 | fn walk(st: &Arc<SearchState>, q: &str, want_name: bool, want_content: bool) { |
| 209 | let abs = st.server_root.join(&st.root.path); |
| 210 | if !abs.is_dir() { |
| 211 | // Root removed out from under us: nothing to search. |
| 212 | return; |
| 213 | } |
| 214 | WalkBuilder::new(&abs) |
| 215 | .standard_filters(true) // fd's defaults: hidden files + gitignore |
| 216 | .build_parallel() |
| 217 | .run(|| { |
| 218 | let abs = abs.clone(); |
| 219 | Box::new(move |result| match result { |
| 220 | Ok(entry) => visit(st, q, want_name, want_content, &entry, &abs), |
| 221 | Err(_) => WalkState::Continue, // unreadable entry: skip like fd |
| 222 | }) |
| 223 | }); |
| 224 | } |
| 225 | |
| 226 | /// One walked entry: emit a name hit, grep it, or both. |
| 227 | fn visit( |
| 228 | st: &Arc<SearchState>, |
| 229 | q: &str, |
| 230 | want_name: bool, |
| 231 | want_content: bool, |
| 232 | entry: &DirEntry, |
| 233 | abs: &Path, |
| 234 | ) -> WalkState { |
| 235 | if st.tx.is_closed() { |
| 236 | return WalkState::Quit; |
| 237 | } |
| 238 | st.scanned.fetch_add(1, Ordering::Relaxed); |
| 239 | // The root itself has an empty relative path and is not a result. |
| 240 | let Some(rel) = entry |
| 241 | .path() |
| 242 | .strip_prefix(abs) |
| 243 | .ok() |
| 244 | .filter(|p| !p.as_os_str().is_empty()) |
| 245 | else { |
| 246 | return WalkState::Continue; |
| 247 | }; |
| 248 | // Original case, unlike the name used for matching: this path is what |
| 249 | // the client opens the entry by. |
| 250 | let rel = rel.to_string_lossy().replace('\\', "/"); |
| 251 | let is_dir = entry.file_type().is_some_and(|t| t.is_dir()); |
| 252 | |
| 253 | if want_name { |
| 254 | // Matched against the entry's own name, not its path: matching the |
| 255 | // path makes every descendant of a matching directory a hit too |
| 256 | // ("e" matching `search-test/` dragged in all 400 files under it), |
| 257 | // which buries the entries the user actually named. |
| 258 | let name = entry.file_name().to_string_lossy().to_lowercase(); |
| 259 | if name_matches(&st.words, &name) { |
| 260 | st.names.fetch_add(1, Ordering::Relaxed); |
| 261 | let size = if is_dir { |
| 262 | 0 |
| 263 | } else { |
| 264 | entry.metadata().map(|m| m.len()).unwrap_or(0) |
| 265 | }; |
| 266 | let ev = SearchEvent::File { |
| 267 | root_id: st.root.id, |
| 268 | path: rel.clone(), |
| 269 | size, |
| 270 | is_dir, |
| 271 | // Same sniff a directory listing does, so the client needs no |
| 272 | // extension table of its own. One open() per name hit, on the |
| 273 | // walker thread that already stat'ed the entry. |
| 274 | kind: crate::fs::detect_kind(entry.path(), is_dir), |
| 275 | }; |
| 276 | // blocking_send: the channel cap provides backpressure against a |
| 277 | // slow client. Err means the receiver is gone: stop the search. |
| 278 | if st.tx.blocking_send(ev).is_err() { |
| 279 | return WalkState::Quit; |
| 280 | } |
| 281 | } |
| 282 | } |
| 283 | |
| 284 | if want_content && entry.file_type().is_some_and(|t| t.is_file()) { |
| 285 | if entry.metadata().map(|m| m.len()).unwrap_or(0) > SEARCH_MAX_FILE_BYTES { |
| 286 | st.skipped.fetch_add(1, Ordering::Relaxed); |
| 287 | return WalkState::Continue; |
| 288 | } |
| 289 | search_one_file(q, entry.path(), rel, st); |
| 290 | if st.tx.is_closed() || st.matches.load(Ordering::Relaxed) >= SEARCH_MAX_MATCH_EVENTS { |
| 291 | return WalkState::Quit; |
| 292 | } |
| 293 | } |
| 294 | WalkState::Continue |
| 295 | } |
| 296 | |
| 297 | /// Grep one file with a fresh literal matcher. The matcher is built per file |
| 298 | /// rather than shared: `RegexMatcher` is not `Sync`, so it cannot cross the |
| 299 | /// walk's thread boundary; a case-insensitive literal compiles in |
| 300 | /// microseconds, negligible next to the file read it protects. |
| 301 | fn search_one_file(q: &str, path: &Path, rel: String, st: &SearchState) { |
| 302 | let matcher = match RegexMatcherBuilder::new() |
| 303 | .case_insensitive(true) |
| 304 | .build_literals(&[q]) |
| 305 | { |
| 306 | Ok(m) => m, |
| 307 | Err(_) => return, |
| 308 | }; |
| 309 | let mut searcher = SearcherBuilder::new().line_number(true).build(); |
| 310 | // ripgrep's default: a NUL byte means "binary, stop". |
| 311 | searcher.set_binary_detection(BinaryDetection::quit(0)); |
| 312 | |
| 313 | let mut sink = MatchSink { st, lines: 0, rel }; |
| 314 | // I/O errors (permissions, vanished file): skip like fd does. |
| 315 | let _ = searcher.search_path(matcher, path, &mut sink); |
| 316 | } |
| 317 | |
| 318 | /// Pushes each matched line onto the event channel. `blocking_send` is |
| 319 | /// correct here: it runs on the walk's worker threads, and the channel cap |
| 320 | /// provides backpressure against a slow client. |
| 321 | struct MatchSink<'a> { |
| 322 | st: &'a SearchState, |
| 323 | /// Matched lines emitted for the current file. |
| 324 | lines: usize, |
| 325 | rel: String, |
| 326 | } |
| 327 | |
| 328 | impl Sink for MatchSink<'_> { |
| 329 | type Error = std::io::Error; |
| 330 | |
| 331 | fn matched(&mut self, _searcher: &Searcher, mat: &SinkMatch) -> Result<bool, Self::Error> { |
| 332 | if self.st.tx.is_closed() || self.lines >= SEARCH_MAX_LINES_PER_FILE { |
| 333 | return Ok(false); |
| 334 | } |
| 335 | let line = match mat.line_number() { |
| 336 | Some(n) => n, |
| 337 | None => return Ok(true), |
| 338 | }; |
| 339 | // Single-line matcher: exactly one line per match. |
| 340 | let text = match mat.lines().next() { |
| 341 | Some(l) => truncate_line(l), |
| 342 | None => return Ok(true), |
| 343 | }; |
| 344 | // Cap check before the increment, so the counted number of matches |
| 345 | // is exactly the number of matches delivered. |
| 346 | if self.st.matches.load(Ordering::Relaxed) >= SEARCH_MAX_MATCH_EVENTS { |
| 347 | return Ok(false); |
| 348 | } |
| 349 | self.lines += 1; |
| 350 | self.st.matches.fetch_add(1, Ordering::Relaxed); |
| 351 | match self.st.tx.blocking_send(SearchEvent::Match { |
| 352 | root_id: self.st.root.id, |
| 353 | path: self.rel.clone(), |
| 354 | line, |
| 355 | text, |
| 356 | }) { |
| 357 | Ok(()) => Ok(true), |
| 358 | // Receiver gone: stop this file; the walker quits on the next |
| 359 | // `is_closed` check. |
| 360 | Err(_) => Ok(false), |
| 361 | } |
| 362 | } |
| 363 | } |
| 364 | |
| 365 | /// Strips the line terminator and truncates at [`SEARCH_MAX_LINE_BYTES`]. |
| 366 | /// |
| 367 | /// A character split by the cap decodes to one replacement character, as |
| 368 | /// does any invalid byte; the cap is therefore a byte cap on the input, not |
| 369 | /// exactly on the output. |
| 370 | fn truncate_line(line: &[u8]) -> String { |
| 371 | let end = line |
| 372 | .iter() |
| 373 | .rposition(|b| *b != b'\n' && *b != b'\r') |
| 374 | .map_or(0, |i| i + 1); |
| 375 | String::from_utf8_lossy(&line[..end.min(SEARCH_MAX_LINE_BYTES)]).into_owned() |
| 376 | } |
| 377 | |
| 378 | #[cfg(test)] |
| 379 | mod tests { |
| 380 | use super::*; |
| 381 | |
| 382 | #[test] |
| 383 | fn name_match_needs_every_word() { |
| 384 | let words = vec!["invoice".to_string(), "2025".to_string()]; |
| 385 | assert!(name_matches(&words, "invoice_2025-11_final.pdf")); |
| 386 | assert!(!name_matches(&words, "invoice_2024.pdf")); |
| 387 | } |
| 388 | |
| 389 | #[test] |
| 390 | fn name_match_ignores_the_parent_path() { |
| 391 | // The caller passes the entry's own name, so a file does not match |
| 392 | // just because an ancestor directory does. |
| 393 | let words = vec!["test".to_string()]; |
| 394 | assert!(name_matches(&words, "test-notes.md")); |
| 395 | assert!(!name_matches(&words, "f12.txt")); |
| 396 | } |
| 397 | |
| 398 | #[test] |
| 399 | fn truncation_strips_terminator() { |
| 400 | assert_eq!(truncate_line(b"hello\n"), "hello"); |
| 401 | assert_eq!(truncate_line(b"hello\r\n"), "hello"); |
| 402 | assert_eq!(truncate_line(b"no newline"), "no newline"); |
| 403 | assert_eq!(truncate_line(b"\n"), ""); |
| 404 | } |
| 405 | |
| 406 | #[test] |
| 407 | fn truncation_caps_length() { |
| 408 | let long = vec![b'x'; 600]; |
| 409 | let t = truncate_line(&long); |
| 410 | assert_eq!(t.len(), SEARCH_MAX_LINE_BYTES); |
| 411 | } |
| 412 | |
| 413 | #[test] |
| 414 | fn truncation_keeps_utf8_boundary() { |
| 415 | // 3-byte chars; the cap must not split one. |
| 416 | let long: Vec<u8> = "ä".repeat(200).into_bytes(); |
| 417 | let t = truncate_line(&long); |
| 418 | assert!(t.len() <= SEARCH_MAX_LINE_BYTES); |
| 419 | assert!(t.chars().all(|c| c == 'ä')); |
| 420 | } |
| 421 | } |
| 422 |