search.rs
⎇
Raw
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
27use std::path::{Path, PathBuf};
28use std::sync::Arc;
29use std::sync::atomic::{AtomicUsize, Ordering};
30use std::time::Instant;
31
32use api_types::SearchEvent;
33use axum::extract::{Query as AxumQuery, State};
34use axum::http::StatusCode;
35use axum::response::sse::{Event, Sse};
36use axum::response::{IntoResponse, Response};
37use futures_util::StreamExt;
38use grep::regex::RegexMatcherBuilder;
39use grep::searcher::{BinaryDetection, Searcher, SearcherBuilder, Sink, SinkMatch};
40use ignore::{DirEntry, WalkBuilder, WalkState};
41use serde::Deserialize;
42
43use crate::api::common::{AuthUser, blocking};
44use crate::db::RootRow;
45use crate::error::{ApiError, AppState};
46
47/// Concurrent searches. A search owns a pool of walker threads, so the two
48/// caps together bound the threads a burst of searches can create.
49static SEARCH_SLOTS: std::sync::LazyLock<Arc<tokio::sync::Semaphore>> =
50 std::sync::LazyLock::new(|| Arc::new(tokio::sync::Semaphore::new(4)));
51
52/// Walker threads per search: the machine's parallelism, capped at 8. The
53/// walk is IO-bound, so more threads buy nothing past that.
54fn walker_threads() -> usize {
55 std::thread::available_parallelism()
56 .map(|n| n.get())
57 .unwrap_or(4)
58 .min(8)
59}
60
61/// Content search skips files larger than this. There is no search budget
62/// (the client can stop at any time), but a multi-gigabyte text file would
63/// stall its worker indefinitely — ripgrep has no such cap, we do.
64const SEARCH_MAX_FILE_BYTES: u64 = 10 * 1024 * 1024;
65/// A matched line is truncated to this many bytes on the wire.
66const SEARCH_MAX_LINE_BYTES: usize = 500;
67/// Total matched lines emitted per search. A one-character query can match
68/// tens of thousands of lines; unbounded streaming would flood the client.
69const SEARCH_MAX_MATCH_EVENTS: usize = 50_000;
70/// Matched lines emitted per file (mirrors the client's render cap).
71const SEARCH_MAX_LINES_PER_FILE: usize = 500;
72
73/// The field names are the shared [`api_types::P_Q`] / [`api_types::P_SCOPE`]
74/// / [`api_types::P_ROOT`] / [`api_types::P_PATH`] constants. `#[serde(rename)]`
75/// only takes a literal, so that link cannot be written here;
76/// `tests::query_fields_are_the_shared_constants` pins it instead.
77#[derive(Debug, Deserialize)]
78pub(super) struct SearchQuery {
79 q: Option<String>,
80 /// `name` (default), `content` or `both`.
81 scope: Option<String>,
82 /// The root to search; omitted = the caller's first root.
83 root: Option<i64>,
84 /// Folder inside the root to start in; omitted = the whole root.
85 path: Option<String>,
86}
87
88pub(super) async fn search(
89 State(state): State<Arc<AppState>>,
90 auth: AuthUser,
91 AxumQuery(query): AxumQuery<SearchQuery>,
92) -> Result<Response, ApiError> {
93 // Search is a session-user feature: a share visitor has no search UI and
94 // must not probe files outside the shared item.
95 if auth.share.is_some() {
96 return Err(ApiError::localized(
97 StatusCode::FORBIDDEN,
98 "search requires a signed-in session",
99 "err_search_forbidden",
100 ));
101 }
102 let q = query
103 .q
104 .as_deref()
105 .map(str::trim)
106 .filter(|q| !q.is_empty())
107 .ok_or_else(|| {
108 ApiError::localized(
109 StatusCode::BAD_REQUEST,
110 "search requires a non-empty query",
111 "err_search_empty_query",
112 )
113 })?;
114
115 let scope = query.scope.as_deref().unwrap_or("name");
116 let want_name = matches!(scope, "name" | "both");
117 let want_content = matches!(scope, "content" | "both");
118 if !(want_name || want_content) {
119 return Err(ApiError::localized(
120 StatusCode::BAD_REQUEST,
121 "scope must be name, content or both",
122 "err_search_bad_scope",
123 ));
124 }
125
126 // Only a root the caller actually has.
127 let root = match query.root {
128 Some(id) => auth.roots.iter().find(|r| r.id == id),
129 None => auth.roots.first(),
130 }
131 .cloned()
132 .ok_or_else(|| {
133 ApiError::localized(
134 StatusCode::FORBIDDEN,
135 "root not accessible",
136 "err_root_forbidden",
137 )
138 })?;
139
140 // The start folder must exist inside the root. Only validated here: the
141 // walk joins the un-canonicalized paths itself, so every result keeps
142 // the root's prefix and stays root-relative.
143 let start_rel = query
144 .path
145 .as_deref()
146 .unwrap_or("")
147 .trim_matches('/')
148 .to_string();
149 let (server_root, root_path, start) =
150 (state.root.clone(), root.path.clone(), start_rel.clone());
151 blocking(move || crate::fs::resolve_dir(&server_root, &root_path, &start)).await?;
152
153 // One slot per running search, held until the walk ends. Each search owns
154 // a pool of walker threads, so unbounded concurrency would swamp the box.
155 // No queueing: a waiting request would hang without any response.
156 let slot = SEARCH_SLOTS.clone().try_acquire_owned().map_err(|_| {
157 ApiError::new(
158 StatusCode::SERVICE_UNAVAILABLE,
159 "too many searches are running",
160 )
161 })?;
162
163 let (tx, mut rx) = tokio::sync::mpsc::channel::<SearchEvent>(256);
164 let search = Arc::new(SearchState {
165 tx,
166 q: q.to_string(),
167 words: q.split_whitespace().map(|w| w.to_lowercase()).collect(),
168 want_name,
169 want_content,
170 root,
171 start_rel,
172 server_root: state.root.clone(),
173 excludes: state.db.search_excludes().await?,
174 started: Instant::now(),
175 scanned: AtomicUsize::new(0),
176 skipped: AtomicUsize::new(0),
177 matches: AtomicUsize::new(0),
178 });
179 run_search(search, slot);
180 let events = futures_util::stream::poll_fn(move |cx| rx.poll_recv(cx));
181 // `Sse` does the `data: <json>\n\n` framing and the content-type and
182 // cache headers; nginx and friends still need telling not to buffer.
183 let sse = Sse::new(events.map(|ev| Event::default().json_data(&ev)));
184 Ok(([("x-accel-buffering", "no")], sse).into_response())
185}
186
187/// Shared state of the search, behind one `Arc` that every walk visitor
188/// clones a handle to. The counters are atomics so they can be updated from
189/// the walker threads lock-free.
190struct SearchState {
191 tx: tokio::sync::mpsc::Sender<SearchEvent>,
192 q: String,
193 words: Vec<String>,
194 want_name: bool,
195 want_content: bool,
196 root: RootRow,
197 /// Where the walk starts, relative to the root ("" = the root itself).
198 start_rel: String,
199 server_root: PathBuf,
200 /// Admin-configured folders left out of every search, relative to the
201 /// server root and already normalised (no slashes at either end).
202 excludes: Vec<String>,
203 started: Instant,
204 /// Directory entries visited.
205 scanned: AtomicUsize,
206 /// Files skipped for content search (over the size cap).
207 skipped: AtomicUsize,
208 matches: AtomicUsize,
209}
210
211/// The search itself: one walk over the root, matching names and grepping
212/// contents as it goes, sending events as they are found. Runs on a plain
213/// thread — `WalkParallel` manages its own worker pool — while the SSE body
214/// stream polls the receiver.
215///
216/// Stopping: when the client goes away, axum drops the body stream, which
217/// drops the receiver. Every `is_closed` check in the walkers then yields
218/// `WalkState::Quit` at the next entry, so the walk unwinds promptly.
219fn run_search(state: Arc<SearchState>, slot: tokio::sync::OwnedSemaphorePermit) {
220 std::thread::spawn(move || {
221 // Released when the walk is done, not when the client stops reading.
222 let _slot = slot;
223 walk(&state);
224 // `stopped`: the receiver went away (client stopped or navigated)
225 // or the match cap was hit, before the walk finished.
226 let stopped = state.tx.is_closed()
227 || state.matches.load(Ordering::Relaxed) >= SEARCH_MAX_MATCH_EVENTS;
228 let _ = state.tx.blocking_send(SearchEvent::Done {
229 stopped,
230 scanned: state.scanned.load(Ordering::Relaxed),
231 skipped: state.skipped.load(Ordering::Relaxed),
232 elapsed_ms: state.started.elapsed().as_millis() as u64,
233 });
234 });
235}
236
237/// A root-relative path re-expressed relative to the server root, the form
238/// the exclude list is stored in. `root_path` is "." for the whole root.
239fn join_rel(root_path: &str, rel: &str) -> String {
240 let root_path = root_path.trim_matches('/');
241 if root_path.is_empty() || root_path == "." {
242 rel.to_string()
243 } else {
244 format!("{root_path}/{rel}")
245 }
246}
247
248/// Whether `path` is an excluded folder or sits under one.
249///
250/// `Path::starts_with` compares whole components, so "docs" does not exclude
251/// the sibling "docs-archive". A plain string prefix would.
252fn is_excluded(excludes: &[String], path: &str) -> bool {
253 let path = Path::new(path);
254 excludes
255 .iter()
256 .any(|e| crate::fs::is_within_or_eq(Path::new(e), path))
257}
258
259/// True when every query word occurs in `name`. Both are already lowercased.
260///
261/// `name` is the entry's own name, never its path — see the module docs.
262fn name_matches(words: &[String], name: &str) -> bool {
263 words.iter().all(|w| name.contains(w.as_str()))
264}
265
266/// Walks the root in parallel, matching each entry against the wanted
267/// scopes. The visitor is cloned once per worker thread, and returns
268/// [`WalkState::Quit`] to stop the whole walk (client went away).
269fn walk(st: &Arc<SearchState>) {
270 let abs = st.server_root.join(&st.root.path);
271 let start = abs.join(&st.start_rel);
272 if !start.is_dir() {
273 // Root or start folder removed out from under us: nothing to search.
274 return;
275 }
276 WalkBuilder::new(&start)
277 // A file browser shows everything, so search must too: no hidden
278 // or `.gitignore`/`.ignore` filtering. Symlinks stay unfollowed.
279 .standard_filters(false)
280 .threads(walker_threads())
281 .build_parallel()
282 .run(|| {
283 let abs = abs.clone();
284 Box::new(move |result| match result {
285 Ok(entry) => visit(st, &entry, &abs),
286 Err(_) => WalkState::Continue, // unreadable entry: skip like fd
287 })
288 });
289}
290
291/// One walked entry: emit a name hit, grep it, or both.
292fn visit(st: &Arc<SearchState>, entry: &DirEntry, abs: &Path) -> WalkState {
293 if st.tx.is_closed() {
294 return WalkState::Quit;
295 }
296 st.scanned.fetch_add(1, Ordering::Relaxed);
297 // The start folder itself (depth 0) is not a result. Paths stay
298 // relative to the root, not to the start folder, so the client opens
299 // them the same way as any listing entry.
300 if entry.depth() == 0 {
301 return WalkState::Continue;
302 }
303 let Some(rel) = entry.path().strip_prefix(abs).ok() else {
304 return WalkState::Continue;
305 };
306 // Original case, unlike the name used for matching: this path is what
307 // the client opens the entry by.
308 let rel = rel.to_string_lossy().replace('\\', "/");
309 let is_dir = entry.file_type().is_some_and(|t| t.is_dir());
310
311 // `Skip` on the folder itself stops the walker descending, so nothing
312 // underneath is ever read.
313 if !st.excludes.is_empty() {
314 let from_server_root = join_rel(&st.root.path, &rel);
315 if is_excluded(&st.excludes, &from_server_root) {
316 return if is_dir {
317 WalkState::Skip
318 } else {
319 WalkState::Continue
320 };
321 }
322 }
323
324 if st.want_name {
325 // Matched against the entry's own name, not its path: matching the
326 // path makes every descendant of a matching directory a hit too
327 // ("e" matching `search-test/` dragged in all 400 files under it),
328 // which buries the entries the user actually named.
329 let name = entry.file_name().to_string_lossy().to_lowercase();
330 if name_matches(&st.words, &name) {
331 let size = if is_dir {
332 0
333 } else {
334 entry.metadata().map(|m| m.len()).unwrap_or(0)
335 };
336 let ev = SearchEvent::File {
337 root_id: st.root.id,
338 path: rel.clone(),
339 size,
340 is_dir,
341 // Same sniff a directory listing does, so the client needs no
342 // extension table of its own. One open() per name hit, on the
343 // walker thread that already stat'ed the entry.
344 kind: crate::fs::detect_kind(entry.path(), is_dir),
345 };
346 // blocking_send: the channel cap provides backpressure against a
347 // slow client. Err means the receiver is gone: stop the search.
348 if st.tx.blocking_send(ev).is_err() {
349 return WalkState::Quit;
350 }
351 }
352 }
353
354 if st.want_content && entry.file_type().is_some_and(|t| t.is_file()) {
355 if entry.metadata().map(|m| m.len()).unwrap_or(0) > SEARCH_MAX_FILE_BYTES {
356 st.skipped.fetch_add(1, Ordering::Relaxed);
357 return WalkState::Continue;
358 }
359 search_one_file(entry.path(), rel, st);
360 if st.tx.is_closed() || st.matches.load(Ordering::Relaxed) >= SEARCH_MAX_MATCH_EVENTS {
361 return WalkState::Quit;
362 }
363 }
364 WalkState::Continue
365}
366
367/// Grep one file with a fresh literal matcher. The matcher is built per file
368/// rather than shared: `RegexMatcher` is not `Sync`, so it cannot cross the
369/// walk's thread boundary; a case-insensitive literal compiles in
370/// microseconds, negligible next to the file read it protects.
371fn search_one_file(path: &Path, rel: String, st: &SearchState) {
372 let matcher = match RegexMatcherBuilder::new()
373 .case_insensitive(true)
374 .build_literals(&[&st.q])
375 {
376 Ok(m) => m,
377 Err(_) => return,
378 };
379 let mut searcher = SearcherBuilder::new().line_number(true).build();
380 // ripgrep's default: a NUL byte means "binary, stop".
381 searcher.set_binary_detection(BinaryDetection::quit(0));
382
383 let mut sink = MatchSink { st, lines: 0, rel };
384 // I/O errors (permissions, vanished file): skip like fd does.
385 let _ = searcher.search_path(matcher, path, &mut sink);
386}
387
388/// Pushes each matched line onto the event channel. `blocking_send` is
389/// correct here: it runs on the walk's worker threads, and the channel cap
390/// provides backpressure against a slow client.
391struct MatchSink<'a> {
392 st: &'a SearchState,
393 /// Matched lines emitted for the current file.
394 lines: usize,
395 rel: String,
396}
397
398impl Sink for MatchSink<'_> {
399 type Error = std::io::Error;
400
401 fn matched(&mut self, _searcher: &Searcher, mat: &SinkMatch) -> Result<bool, Self::Error> {
402 if self.st.tx.is_closed() || self.lines >= SEARCH_MAX_LINES_PER_FILE {
403 return Ok(false);
404 }
405 let line = match mat.line_number() {
406 Some(n) => n,
407 None => return Ok(true),
408 };
409 // Single-line matcher: exactly one line per match.
410 let text = match mat.lines().next() {
411 Some(l) => truncate_line(l),
412 None => return Ok(true),
413 };
414 // Cap check before the increment, so the counted number of matches
415 // is exactly the number of matches delivered.
416 if self.st.matches.load(Ordering::Relaxed) >= SEARCH_MAX_MATCH_EVENTS {
417 return Ok(false);
418 }
419 self.lines += 1;
420 self.st.matches.fetch_add(1, Ordering::Relaxed);
421 match self.st.tx.blocking_send(SearchEvent::Match {
422 root_id: self.st.root.id,
423 path: self.rel.clone(),
424 line,
425 text,
426 }) {
427 Ok(()) => Ok(true),
428 // Receiver gone: stop this file; the walker quits on the next
429 // `is_closed` check.
430 Err(_) => Ok(false),
431 }
432 }
433}
434
435/// Strips the line terminator and truncates at [`SEARCH_MAX_LINE_BYTES`].
436///
437/// A character split by the cap decodes to one replacement character, as
438/// does any invalid byte; the cap is therefore a byte cap on the input, not
439/// exactly on the output.
440fn truncate_line(line: &[u8]) -> String {
441 let end = line
442 .iter()
443 .rposition(|b| *b != b'\n' && *b != b'\r')
444 .map_or(0, |i| i + 1);
445 String::from_utf8_lossy(&line[..end.min(SEARCH_MAX_LINE_BYTES)]).into_owned()
446}
447
448#[cfg(test)]
449mod tests {
450 use super::*;
451
452 /// A sibling whose name merely starts with an excluded folder's name
453 /// must still be searchable.
454 #[test]
455 fn excludes_cover_descendants_but_not_name_siblings() {
456 let ex = vec!["private".to_string(), "a/b".to_string()];
457 assert!(is_excluded(&ex, "private"));
458 assert!(is_excluded(&ex, "private/deep/file.txt"));
459 assert!(is_excluded(&ex, "a/b"));
460 assert!(is_excluded(&ex, "a/b/c.txt"));
461
462 assert!(!is_excluded(&ex, "private-archive"));
463 assert!(!is_excluded(&ex, "privateer.txt"));
464 assert!(!is_excluded(&ex, "a"));
465 assert!(!is_excluded(&ex, "a/bc"));
466 assert!(!is_excluded(&ex, "other/private"));
467 assert!(!is_excluded(&[], "anything"));
468 }
469
470 /// Results are root-relative; the exclude list is server-root-relative.
471 #[test]
472 fn join_rel_lifts_a_path_to_the_server_root() {
473 assert_eq!(join_rel(".", "docs/a.txt"), "docs/a.txt");
474 assert_eq!(join_rel("", "docs/a.txt"), "docs/a.txt");
475 assert_eq!(join_rel("home/bob", "docs/a.txt"), "home/bob/docs/a.txt");
476 assert_eq!(join_rel("/home/bob/", "x"), "home/bob/x");
477 }
478
479 /// A rename of one of the `P_*` constants without the matching field
480 /// rename would silently stop the server from reading the parameter the
481 /// client sends. This builds the query string from the constants and
482 /// runs the real extractor over it.
483 #[test]
484 fn query_fields_are_the_shared_constants() {
485 use api_types::{P_PATH, P_Q, P_ROOT, P_SCOPE};
486 let uri: axum::http::Uri =
487 format!("/api/search?{P_Q}=invoice&{P_SCOPE}=both&{P_ROOT}=7&{P_PATH}=docs")
488 .parse()
489 .unwrap();
490 let q: SearchQuery = AxumQuery::try_from_uri(&uri).unwrap().0;
491 assert_eq!(q.q.as_deref(), Some("invoice"));
492 assert_eq!(q.scope.as_deref(), Some("both"));
493 assert_eq!(q.root, Some(7));
494 assert_eq!(q.path.as_deref(), Some("docs"));
495 }
496
497 #[test]
498 fn name_match_needs_every_word() {
499 let words = vec!["invoice".to_string(), "2025".to_string()];
500 assert!(name_matches(&words, "invoice_2025-11_final.pdf"));
501 assert!(!name_matches(&words, "invoice_2024.pdf"));
502 }
503
504 #[test]
505 fn name_match_ignores_the_parent_path() {
506 // The caller passes the entry's own name, so a file does not match
507 // just because an ancestor directory does.
508 let words = vec!["test".to_string()];
509 assert!(name_matches(&words, "test-notes.md"));
510 assert!(!name_matches(&words, "f12.txt"));
511 }
512
513 #[test]
514 fn truncation_strips_terminator() {
515 assert_eq!(truncate_line(b"hello\n"), "hello");
516 assert_eq!(truncate_line(b"hello\r\n"), "hello");
517 assert_eq!(truncate_line(b"no newline"), "no newline");
518 assert_eq!(truncate_line(b"\n"), "");
519 }
520
521 #[test]
522 fn truncation_caps_length() {
523 let long = vec![b'x'; 600];
524 let t = truncate_line(&long);
525 assert_eq!(t.len(), SEARCH_MAX_LINE_BYTES);
526 }
527
528 #[test]
529 fn truncation_keeps_utf8_boundary() {
530 // 3-byte chars; the cap must not split one.
531 let long: Vec<u8> = "ä".repeat(200).into_bytes();
532 let t = truncate_line(&long);
533 assert!(t.len() <= SEARCH_MAX_LINE_BYTES);
534 assert!(t.chars().all(|c| c == 'ä'));
535 }
536}
537