//! `GET /api/search` — index-free name and content search, streamed as SSE. //! //! No index, by design: every search walks the selected roots with //! [`ignore::WalkParallel`] (fd/ripgrep's walker: parallel, skips hidden //! files, honors `.gitignore`) 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. //! //! 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::http::header; use axum::http::request::Parts; use axum::response::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, HasState}; use crate::db::RootRow; use crate::error::{ApiError, AppState}; /// 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; #[derive(Debug, Deserialize)] pub(super) struct SearchQuery { q: Option, /// `name` (default), `content` or `both`. #[serde(default)] scope: Option, /// Comma-separated root ids; empty = all of the user's roots. #[serde(default)] roots: Option, } /// Search is a session-user feature: reject share-token sessions outright /// (the share page has no search UI, and a share visitor must not probe /// files outside the shared item). pub(super) struct SearchUser(AuthUser); impl axum::extract::FromRequestParts for SearchUser where S: HasState + Send + Sync, { type Rejection = ApiError; async fn from_request_parts(parts: &mut Parts, state: &S) -> Result { let auth = AuthUser::from_request_parts(parts, state).await?; if auth.share.is_some() { return Err(ApiError::localized( StatusCode::FORBIDDEN, "search requires a signed-in session", "err_search_forbidden", )); } Ok(SearchUser(auth)) } } pub(super) async fn search( State(state): State>, SearchUser(auth): SearchUser, AxumQuery(query): AxumQuery, ) -> Result { 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", )); } // Selected roots: only ids the caller actually has. let roots = match query.roots.as_deref() { None => Ok(auth.roots.clone()), Some(list) if list.trim().is_empty() => Ok(auth.roots.clone()), Some(list) => list .split(',') .map(str::trim) .filter(|s| !s.is_empty()) .map(|s| -> Result { let id = s.parse::().map_err(|_| { ApiError::localized( StatusCode::BAD_REQUEST, "root ids must be numbers", "err_search_bad_roots", ) })?; auth.roots .iter() .find(|r| r.id == id) .cloned() .ok_or_else(|| { ApiError::localized( StatusCode::FORBIDDEN, "root not accessible", "err_root_forbidden", ) }) }) .collect::, ApiError>>(), }?; Ok(sse_response(search_stream( q.to_string(), want_name, want_content, Arc::new(roots), state.root.clone(), ))) } /// Wraps an event stream in an SSE response (`data: ` per event). fn sse_response( events: impl futures_util::Stream + Send + 'static, ) -> Response { let body = events.map(|ev| { // Infallible: SearchEvent is always serializable. let json = serde_json::to_string(&ev).expect("SearchEvent serializes"); Ok::(format!("data: {json}\n\n")) }); Response::builder() .status(StatusCode::OK) .header(header::CONTENT_TYPE, "text/event-stream; charset=utf-8") .header(header::CACHE_CONTROL, "no-cache") // nginx and friends must not buffer an event stream. .header("x-accel-buffering", "no") .body(axum::body::Body::from_stream(body)) .expect("static response") } /// Shared state of the search. The counters are `Arc`'d so the per-thread /// walk visitors (which must be `Clone`) can update them lock-free. struct SearchState { tx: tokio::sync::mpsc::Sender, words: Vec, roots: Arc>, server_root: PathBuf, started: Instant, /// Directory entries visited (both phases). scanned: Arc, /// Files skipped by the content phase (over the size cap). skipped: Arc, names: Arc, matches: Arc, } /// The search itself: walk (name phase) and grep (content phase) over the /// selected roots, sending events as they are found. Runs on a plain thread /// — each `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. fn search_stream( q: String, want_name: bool, want_content: bool, roots: Arc>, server_root: PathBuf, ) -> 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(), roots, server_root, started: Instant::now(), scanned: Arc::new(AtomicUsize::new(0)), skipped: Arc::new(AtomicUsize::new(0)), names: Arc::new(AtomicUsize::new(0)), matches: Arc::new(AtomicUsize::new(0)), }); std::thread::spawn(move || { if want_name { name_phase(&state); } if want_content { // The name phase already counted every visited entry; only // count again when the content phase is the only walk. content_phase(&state, &q, !want_name); } // `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) } /// Walks one root in parallel. `visit` sees every entry (the root itself /// included) plus the root's absolute path, and returns [`WalkState::Quit`] /// to stop the whole walk (client went away). Must be `Clone`: the walker /// clones it once per worker thread. fn walk_root(root: &RootRow, server_root: &Path, visit: F) where F: FnMut(&DirEntry, &Path) -> WalkState + Clone + Send, { let abs = server_root.join(&root.path); if !abs.is_dir() { // Root removed out from under us: skip, keep searching the rest. return; } WalkBuilder::new(&abs) .standard_filters(true) // fd's defaults: hidden files + gitignore .build_parallel() .run(move || { let mut visit = visit.clone(); let abs = abs.clone(); Box::new(move |result| match result { Ok(entry) => visit(&entry, &abs), Err(_) => WalkState::Continue, // unreadable entry: skip like fd }) }); } /// 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())) } /// Phase 1: name matches. fn name_phase(st: &Arc) { for root in st.roots.iter() { walk_root(root, &st.server_root, move |entry, abs| { if st.tx.is_closed() { return WalkState::Quit; } st.scanned.fetch_add(1, Ordering::Relaxed); let Some(rel) = entry .path() .strip_prefix(abs) .ok() .filter(|p| !p.as_os_str().is_empty()) else { return WalkState::Continue; }; // 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) { 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('\\', "/"); st.names.fetch_add(1, Ordering::Relaxed); let is_dir = entry.file_type().map(|t| t.is_dir()).unwrap_or(false); let size = if is_dir { 0 } else { entry.metadata().map(|m| m.len()).unwrap_or(0) }; let ev = SearchEvent::File { root_id: root.id, path: rel, size, is_dir, }; // blocking_send: the channel cap provides backpressure against a // slow client. Err means the receiver is gone: stop the search. match st.tx.blocking_send(ev) { Ok(()) => WalkState::Continue, Err(_) => WalkState::Quit, } }); } } /// Phase 2: content matches. fn content_phase(st: &Arc, q: &str, count_scanned: bool) { for root in st.roots.iter() { walk_root(root, &st.server_root, move |entry, abs| { if st.tx.is_closed() { return WalkState::Quit; } if count_scanned { st.scanned.fetch_add(1, Ordering::Relaxed); } if !entry.file_type().map(|t| t.is_file()).unwrap_or(false) { return WalkState::Continue; } if entry.metadata().map(|m| m.len()).unwrap_or(0) > SEARCH_MAX_FILE_BYTES { st.skipped.fetch_add(1, Ordering::Relaxed); return WalkState::Continue; } let Some(rel) = entry.path().strip_prefix(abs).ok() else { return WalkState::Continue; }; search_one_file(q, entry.path(), rel, root.id, &st.tx, &st.matches); if st.tx.is_closed() || st.matches.load(Ordering::Relaxed) >= SEARCH_MAX_MATCH_EVENTS { WalkState::Quit } else { 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: &Path, root_id: i64, tx: &tokio::sync::mpsc::Sender, matches: &Arc, ) { 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 rel = rel.to_string_lossy().replace('\\', "/"); let mut sink = MatchSink { tx, matches, lines: 0, root_id, 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> { tx: &'a tokio::sync::mpsc::Sender, matches: &'a Arc, /// Matched lines emitted for the current file. lines: usize, root_id: i64, rel: String, } impl Sink for MatchSink<'_> { type Error = std::io::Error; fn matched(&mut self, _searcher: &Searcher, mat: &SinkMatch) -> Result { if self.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.matches.load(Ordering::Relaxed) >= SEARCH_MAX_MATCH_EVENTS { return Ok(false); } self.lines += 1; self.matches.fetch_add(1, Ordering::Relaxed); match self.tx.blocking_send(SearchEvent::Match { root_id: self.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`], /// backing off to a UTF-8 boundary if the cut split a character. fn truncate_line(line: &[u8]) -> String { let end = line .iter() .rposition(|b| *b != b'\n' && *b != b'\r') .map(|i| i + 1) .unwrap_or(0); let full = &line[..end]; let cut = full.len().min(SEARCH_MAX_LINE_BYTES); let mut line = &full[..cut]; if cut < full.len() && std::str::from_utf8(full).is_ok() { while !line.is_empty() && std::str::from_utf8(line).is_err() { line = &line[..line.len() - 1]; } } String::from_utf8_lossy(line).into_owned() } #[cfg(test)] mod tests { use super::*; #[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 == 'ä')); } }