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