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