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