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