//! `GET /api/search` — index-free name and content search, streamed as SSE. //! //! No index, by design: every search walks the selected root with //! [`ignore::WalkParallel`] (fd/ripgrep's parallel walker, with all its //! filters off: hidden and gitignored entries are searched too) and matches //! on the fly. //! //! * **Name** (scope `name`/`both`): every word of the query must occur //! (case-insensitive) in the entry's own name — the last path component, //! not the whole relative path — so "invoice 2025" matches //! `invoice_2025-11_final.pdf` but a file is not a hit merely for sitting //! inside a directory whose name matches. //! * **Content** (scope `content`/`both`): the whole query as a literal, //! case-insensitive, ripgrep-style scan via the `grep` crates (ripgrep's //! engine). Binary files are skipped by NUL detection, files over //! [`SEARCH_MAX_FILE_BYTES`] are skipped and counted. //! //! Both scopes run in one walk, so `both` reads the tree once and its file //! and match events interleave. //! //! Results are one JSON object per SSE event, in the order found; the stream //! always ends with a `done` event. The search stops when the client goes //! away: axum drops the response body stream, the channel receiver is //! dropped with it, and the walkers bail at the next `is_closed` check — //! no per-search budgets or timeouts. use std::path::{Path, PathBuf}; use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; use std::time::Instant; use api_types::SearchEvent; use axum::extract::{Query as AxumQuery, State}; use axum::http::StatusCode; use axum::response::sse::{Event, Sse}; use axum::response::{IntoResponse, Response}; use futures_util::StreamExt; use grep::regex::RegexMatcherBuilder; use grep::searcher::{BinaryDetection, Searcher, SearcherBuilder, Sink, SinkMatch}; use ignore::{DirEntry, WalkBuilder, WalkState}; use serde::Deserialize; use crate::api::common::{AuthUser, blocking}; use crate::db::RootRow; use crate::error::{ApiError, AppState}; /// Concurrent searches. A search owns a pool of walker threads, so the two /// caps together bound the threads a burst of searches can create. static SEARCH_SLOTS: std::sync::LazyLock> = std::sync::LazyLock::new(|| Arc::new(tokio::sync::Semaphore::new(4))); /// Walker threads per search: the machine's parallelism, capped at 8. The /// walk is IO-bound, so more threads buy nothing past that. fn walker_threads() -> usize { std::thread::available_parallelism() .map(|n| n.get()) .unwrap_or(4) .min(8) } /// Content search skips files larger than this. There is no search budget /// (the client can stop at any time), but a multi-gigabyte text file would /// stall its worker indefinitely — ripgrep has no such cap, we do. const SEARCH_MAX_FILE_BYTES: u64 = 10 * 1024 * 1024; /// A matched line is truncated to this many bytes on the wire. const SEARCH_MAX_LINE_BYTES: usize = 500; /// Total matched lines emitted per search. A one-character query can match /// tens of thousands of lines; unbounded streaming would flood the client. const SEARCH_MAX_MATCH_EVENTS: usize = 50_000; /// Matched lines emitted per file (mirrors the client's render cap). const SEARCH_MAX_LINES_PER_FILE: usize = 500; /// The field names are the shared [`api_types::P_Q`] / [`api_types::P_SCOPE`] /// / [`api_types::P_ROOT`] / [`api_types::P_PATH`] constants. `#[serde(rename)]` /// only takes a literal, so that link cannot be written here; /// `tests::query_fields_are_the_shared_constants` pins it instead. #[derive(Debug, Deserialize)] pub(super) struct SearchQuery { q: Option, /// `name` (default), `content` or `both`. #[serde(default)] scope: Option, /// The root to search; omitted = the caller's first root. #[serde(default)] root: Option, /// Folder inside the root to start in; omitted = the whole root. #[serde(default)] path: Option, } pub(super) async fn search( State(state): State>, auth: AuthUser, AxumQuery(query): AxumQuery, ) -> Result { // Search is a session-user feature: a share visitor has no search UI and // must not probe files outside the shared item. if auth.share.is_some() { return Err(ApiError::localized( StatusCode::FORBIDDEN, "search requires a signed-in session", "err_search_forbidden", )); } let q = query .q .as_deref() .map(str::trim) .filter(|q| !q.is_empty()) .ok_or_else(|| { ApiError::localized( StatusCode::BAD_REQUEST, "search requires a non-empty query", "err_search_empty_query", ) })?; let scope = query.scope.as_deref().unwrap_or("name"); let want_name = matches!(scope, "name" | "both"); let want_content = matches!(scope, "content" | "both"); if !(want_name || want_content) { return Err(ApiError::localized( StatusCode::BAD_REQUEST, "scope must be name, content or both", "err_search_bad_scope", )); } // Only a root the caller actually has. let root = match query.root { Some(id) => auth.roots.iter().find(|r| r.id == id), None => auth.roots.first(), } .cloned() .ok_or_else(|| { ApiError::localized( StatusCode::FORBIDDEN, "root not accessible", "err_root_forbidden", ) })?; // The start folder must exist inside the root. Only validated here: the // walk joins the un-canonicalized paths itself, so every result keeps // the root's prefix and stays root-relative. let start_rel = query .path .as_deref() .unwrap_or("") .trim_matches('/') .to_string(); let (server_root, root_path, start) = (state.root.clone(), root.path.clone(), start_rel.clone()); blocking(move || crate::fs::resolve_dir(&server_root, &root_path, &start)).await?; // One slot per running search, held until the walk ends. Each search owns // a pool of walker threads, so unbounded concurrency would swamp the box. // No queueing: a waiting request would hang without any response. let slot = SEARCH_SLOTS.clone().try_acquire_owned().map_err(|_| { ApiError::new( StatusCode::SERVICE_UNAVAILABLE, "too many searches are running", ) })?; let events = search_stream( q.to_string(), want_name, want_content, root, start_rel, state.root.clone(), state.db.search_excludes().await?, slot, ); // `Sse` does the `data: \n\n` framing and the content-type and // cache headers; nginx and friends still need telling not to buffer. let sse = Sse::new(events.map(|ev| Event::default().json_data(&ev))); Ok(([("x-accel-buffering", "no")], sse).into_response()) } /// Shared state of the search, behind one `Arc` that every walk visitor /// clones a handle to. The counters are atomics so they can be updated from /// the walker threads lock-free. struct SearchState { tx: tokio::sync::mpsc::Sender, words: Vec, root: RootRow, /// Where the walk starts, relative to the root ("" = the root itself). start_rel: String, server_root: PathBuf, /// Admin-configured folders left out of every search, relative to the /// server root and already normalised (no slashes at either end). excludes: Vec, started: Instant, /// Directory entries visited. scanned: AtomicUsize, /// Files skipped for content search (over the size cap). skipped: AtomicUsize, names: AtomicUsize, matches: AtomicUsize, } /// The search itself: one walk over the root, matching names and grepping /// contents as it goes, sending events as they are found. Runs on a plain /// thread — `WalkParallel` manages its own worker pool — while the SSE body /// stream polls the receiver. /// /// Stopping: when the client goes away, axum drops the body stream, which /// drops the receiver. Every `is_closed` check in the walkers then yields /// `WalkState::Quit` at the next entry, so the walk unwinds promptly. #[allow(clippy::too_many_arguments)] // one search's whole configuration fn search_stream( q: String, want_name: bool, want_content: bool, root: RootRow, start_rel: String, server_root: PathBuf, excludes: Vec, slot: tokio::sync::OwnedSemaphorePermit, ) -> impl futures_util::Stream + Send { let (tx, rx) = tokio::sync::mpsc::channel::(256); let state = Arc::new(SearchState { tx, words: q.split_whitespace().map(|w| w.to_lowercase()).collect(), root, start_rel, server_root, excludes, started: Instant::now(), scanned: AtomicUsize::new(0), skipped: AtomicUsize::new(0), names: AtomicUsize::new(0), matches: AtomicUsize::new(0), }); std::thread::spawn(move || { // Released when the walk is done, not when the client stops reading. let _slot = slot; walk(&state, &q, want_name, want_content); // `stopped`: the receiver went away (client stopped or navigated) // or the match cap was hit, before the walk finished. let stopped = state.tx.is_closed() || state.matches.load(Ordering::Relaxed) >= SEARCH_MAX_MATCH_EVENTS; let _ = state.tx.blocking_send(SearchEvent::Done { stopped, files: state.names.load(Ordering::Relaxed), matches: state.matches.load(Ordering::Relaxed), scanned: state.scanned.load(Ordering::Relaxed), skipped: state.skipped.load(Ordering::Relaxed), elapsed_ms: state.started.elapsed().as_millis() as u64, }); }); tokio_stream::wrappers::ReceiverStream::new(rx) } /// A root-relative path re-expressed relative to the server root, the form /// the exclude list is stored in. `root_path` is "." for the whole root. fn join_rel(root_path: &str, rel: &str) -> String { let root_path = root_path.trim_matches('/'); if root_path.is_empty() || root_path == "." { rel.to_string() } else { format!("{root_path}/{rel}") } } /// Whether `path` is an excluded folder or sits under one. /// /// `Path::starts_with` compares whole components, so "docs" does not exclude /// the sibling "docs-archive". A plain string prefix would. fn is_excluded(excludes: &[String], path: &str) -> bool { let path = Path::new(path); excludes .iter() .any(|e| crate::fs::is_within_or_eq(Path::new(e), path)) } /// True when every query word occurs in `name`. Both are already lowercased. /// /// `name` is the entry's own name, never its path — see the module docs. fn name_matches(words: &[String], name: &str) -> bool { words.iter().all(|w| name.contains(w.as_str())) } /// Walks the root in parallel, matching each entry against the wanted /// scopes. The visitor is cloned once per worker thread, and returns /// [`WalkState::Quit`] to stop the whole walk (client went away). fn walk(st: &Arc, q: &str, want_name: bool, want_content: bool) { let abs = st.server_root.join(&st.root.path); let start = abs.join(&st.start_rel); if !start.is_dir() { // Root or start folder removed out from under us: nothing to search. return; } WalkBuilder::new(&start) // A file browser shows everything, so search must too: no hidden // or `.gitignore`/`.ignore` filtering. Symlinks stay unfollowed. .standard_filters(false) .threads(walker_threads()) .build_parallel() .run(|| { let abs = abs.clone(); Box::new(move |result| match result { Ok(entry) => visit(st, q, want_name, want_content, &entry, &abs), Err(_) => WalkState::Continue, // unreadable entry: skip like fd }) }); } /// One walked entry: emit a name hit, grep it, or both. fn visit( st: &Arc, q: &str, want_name: bool, want_content: bool, entry: &DirEntry, abs: &Path, ) -> WalkState { if st.tx.is_closed() { return WalkState::Quit; } st.scanned.fetch_add(1, Ordering::Relaxed); // The start folder itself (depth 0) is not a result. Paths stay // relative to the root, not to the start folder, so the client opens // them the same way as any listing entry. if entry.depth() == 0 { return WalkState::Continue; } let Some(rel) = entry.path().strip_prefix(abs).ok() else { return WalkState::Continue; }; // Original case, unlike the name used for matching: this path is what // the client opens the entry by. let rel = rel.to_string_lossy().replace('\\', "/"); let is_dir = entry.file_type().is_some_and(|t| t.is_dir()); // `Skip` on the folder itself stops the walker descending, so nothing // underneath is ever read. if !st.excludes.is_empty() { let from_server_root = join_rel(&st.root.path, &rel); if is_excluded(&st.excludes, &from_server_root) { return if is_dir { WalkState::Skip } else { WalkState::Continue }; } } if want_name { // Matched against the entry's own name, not its path: matching the // path makes every descendant of a matching directory a hit too // ("e" matching `search-test/` dragged in all 400 files under it), // which buries the entries the user actually named. let name = entry.file_name().to_string_lossy().to_lowercase(); if name_matches(&st.words, &name) { st.names.fetch_add(1, Ordering::Relaxed); let size = if is_dir { 0 } else { entry.metadata().map(|m| m.len()).unwrap_or(0) }; let ev = SearchEvent::File { root_id: st.root.id, path: rel.clone(), size, is_dir, // Same sniff a directory listing does, so the client needs no // extension table of its own. One open() per name hit, on the // walker thread that already stat'ed the entry. kind: crate::fs::detect_kind(entry.path(), is_dir), }; // blocking_send: the channel cap provides backpressure against a // slow client. Err means the receiver is gone: stop the search. if st.tx.blocking_send(ev).is_err() { return WalkState::Quit; } } } if want_content && entry.file_type().is_some_and(|t| t.is_file()) { if entry.metadata().map(|m| m.len()).unwrap_or(0) > SEARCH_MAX_FILE_BYTES { st.skipped.fetch_add(1, Ordering::Relaxed); return WalkState::Continue; } search_one_file(q, entry.path(), rel, st); if st.tx.is_closed() || st.matches.load(Ordering::Relaxed) >= SEARCH_MAX_MATCH_EVENTS { return WalkState::Quit; } } WalkState::Continue } /// Grep one file with a fresh literal matcher. The matcher is built per file /// rather than shared: `RegexMatcher` is not `Sync`, so it cannot cross the /// walk's thread boundary; a case-insensitive literal compiles in /// microseconds, negligible next to the file read it protects. fn search_one_file(q: &str, path: &Path, rel: String, st: &SearchState) { let matcher = match RegexMatcherBuilder::new() .case_insensitive(true) .build_literals(&[q]) { Ok(m) => m, Err(_) => return, }; let mut searcher = SearcherBuilder::new().line_number(true).build(); // ripgrep's default: a NUL byte means "binary, stop". searcher.set_binary_detection(BinaryDetection::quit(0)); let mut sink = MatchSink { st, lines: 0, rel }; // I/O errors (permissions, vanished file): skip like fd does. let _ = searcher.search_path(matcher, path, &mut sink); } /// Pushes each matched line onto the event channel. `blocking_send` is /// correct here: it runs on the walk's worker threads, and the channel cap /// provides backpressure against a slow client. struct MatchSink<'a> { st: &'a SearchState, /// Matched lines emitted for the current file. lines: usize, rel: String, } impl Sink for MatchSink<'_> { type Error = std::io::Error; fn matched(&mut self, _searcher: &Searcher, mat: &SinkMatch) -> Result { if self.st.tx.is_closed() || self.lines >= SEARCH_MAX_LINES_PER_FILE { return Ok(false); } let line = match mat.line_number() { Some(n) => n, None => return Ok(true), }; // Single-line matcher: exactly one line per match. let text = match mat.lines().next() { Some(l) => truncate_line(l), None => return Ok(true), }; // Cap check before the increment, so the counted number of matches // is exactly the number of matches delivered. if self.st.matches.load(Ordering::Relaxed) >= SEARCH_MAX_MATCH_EVENTS { return Ok(false); } self.lines += 1; self.st.matches.fetch_add(1, Ordering::Relaxed); match self.st.tx.blocking_send(SearchEvent::Match { root_id: self.st.root.id, path: self.rel.clone(), line, text, }) { Ok(()) => Ok(true), // Receiver gone: stop this file; the walker quits on the next // `is_closed` check. Err(_) => Ok(false), } } } /// Strips the line terminator and truncates at [`SEARCH_MAX_LINE_BYTES`]. /// /// A character split by the cap decodes to one replacement character, as /// does any invalid byte; the cap is therefore a byte cap on the input, not /// exactly on the output. fn truncate_line(line: &[u8]) -> String { let end = line .iter() .rposition(|b| *b != b'\n' && *b != b'\r') .map_or(0, |i| i + 1); String::from_utf8_lossy(&line[..end.min(SEARCH_MAX_LINE_BYTES)]).into_owned() } #[cfg(test)] mod tests { use super::*; /// A sibling whose name merely starts with an excluded folder's name /// must still be searchable. #[test] fn excludes_cover_descendants_but_not_name_siblings() { let ex = vec!["private".to_string(), "a/b".to_string()]; assert!(is_excluded(&ex, "private")); assert!(is_excluded(&ex, "private/deep/file.txt")); assert!(is_excluded(&ex, "a/b")); assert!(is_excluded(&ex, "a/b/c.txt")); assert!(!is_excluded(&ex, "private-archive")); assert!(!is_excluded(&ex, "privateer.txt")); assert!(!is_excluded(&ex, "a")); assert!(!is_excluded(&ex, "a/bc")); assert!(!is_excluded(&ex, "other/private")); assert!(!is_excluded(&[], "anything")); } /// Results are root-relative; the exclude list is server-root-relative. #[test] fn join_rel_lifts_a_path_to_the_server_root() { assert_eq!(join_rel(".", "docs/a.txt"), "docs/a.txt"); assert_eq!(join_rel("", "docs/a.txt"), "docs/a.txt"); assert_eq!(join_rel("home/bob", "docs/a.txt"), "home/bob/docs/a.txt"); assert_eq!(join_rel("/home/bob/", "x"), "home/bob/x"); } /// A rename of one of the `P_*` constants without the matching field /// rename would silently stop the server from reading the parameter the /// client sends. This builds the query string from the constants and /// runs the real extractor over it. #[test] fn query_fields_are_the_shared_constants() { use api_types::{P_PATH, P_Q, P_ROOT, P_SCOPE}; let uri: axum::http::Uri = format!("/api/search?{P_Q}=invoice&{P_SCOPE}=both&{P_ROOT}=7&{P_PATH}=docs") .parse() .unwrap(); let q: SearchQuery = AxumQuery::try_from_uri(&uri).unwrap().0; assert_eq!(q.q.as_deref(), Some("invoice")); assert_eq!(q.scope.as_deref(), Some("both")); assert_eq!(q.root, Some(7)); assert_eq!(q.path.as_deref(), Some("docs")); } #[test] fn name_match_needs_every_word() { let words = vec!["invoice".to_string(), "2025".to_string()]; assert!(name_matches(&words, "invoice_2025-11_final.pdf")); assert!(!name_matches(&words, "invoice_2024.pdf")); } #[test] fn name_match_ignores_the_parent_path() { // The caller passes the entry's own name, so a file does not match // just because an ancestor directory does. let words = vec!["test".to_string()]; assert!(name_matches(&words, "test-notes.md")); assert!(!name_matches(&words, "f12.txt")); } #[test] fn truncation_strips_terminator() { assert_eq!(truncate_line(b"hello\n"), "hello"); assert_eq!(truncate_line(b"hello\r\n"), "hello"); assert_eq!(truncate_line(b"no newline"), "no newline"); assert_eq!(truncate_line(b"\n"), ""); } #[test] fn truncation_caps_length() { let long = vec![b'x'; 600]; let t = truncate_line(&long); assert_eq!(t.len(), SEARCH_MAX_LINE_BYTES); } #[test] fn truncation_keeps_utf8_boundary() { // 3-byte chars; the cap must not split one. let long: Vec = "ä".repeat(200).into_bytes(); let t = truncate_line(&long); assert!(t.len() <= SEARCH_MAX_LINE_BYTES); assert!(t.chars().all(|c| c == 'ä')); } }