Search: cut the scaffolding around it
An over-engineering pass over the search feature. Same behaviour, ~580 fewer lines. Server: - one walk serves both scopes instead of two, so `both` reads the tree once; its file and match events now interleave - axum's `Sse` does the framing and the headers; the hand-built `data: <json>` writer and the `SearchUser` extractor are gone - `?roots=` (a comma list the client never sent more than one of) becomes `?root=`, and the counters lose their redundant `Arc` - `SearchEvent::File` carries the sniffed `FileKind`, so the client needs no extension table of its own Client: - `EventSource` replaces ~80 lines of fetch + stream reader + manual SSE framing. It closes on the `error` that ends a stream, so the browser cannot reconnect and re-run the search. The cost is that a rejected request reports one generic message: an `EventSource` cannot read the server's JSON error body. - `UrlSearchParams` replaces the hand-rolled hash-query parser - the chunk-size feedback controller becomes one constant, the `Frag` enum folds into `highlight_html`, and its char-remap for characters whose lowercase expands is gone - one staleness flag instead of two, one `Status::Done` instead of three variants, query words hoisted out of the per-row render Dropped the committed design mockup. Added an integration test for the merged walk, name-vs-path matching, the SSE framing and the rejection paths. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
D.search-mockups/b2-view-streaming.html-222
@@ -1,222 +0,0 @@
<!DOCTYPE html>
<html lang="en" class="dark">
<head>
<meta charset="utf-8">
<title>Mockup B2 — Search view, streaming</title>
<link rel="preconnect" href="https://fonts.googleapis.com">
<link href="https://fonts.googleapis.com/css2?family=Material+Symbols+Rounded:opsz,wght@20..24,400&family=JetBrains+Mono:wght@400;600&display=swap" rel="stylesheet">
<style>
:root {
color-scheme: dark;
--bg: #0f1115;
--panel: #171a20;
--text: #e5e7eb;
--muted: #9ca3af;
--accent: #60a5fa;
--accent-hover: #3b82f6;
--accent-border: #79b5fc;
--danger: #f87171;
--danger-border: #fa8b8b;
--border: #2a2f3a;
--topbar-h: 48px;
--mark: color-mix(in srgb, var(--accent) 35%, transparent);
}
* { box-sizing: border-box; border-radius: 0; }
body { margin: 0; font-family: system-ui, -apple-system, "Segoe UI", Roboto, sans-serif; background: var(--bg); color: var(--text); font-size: 15px; }
.material-symbols-rounded { font-family: 'Material Symbols Rounded'; font-weight: normal; font-style: normal; line-height: 1; letter-spacing: normal; text-transform: none; display: inline-block; white-space: nowrap; direction: ltr; font-variation-settings: 'FILL' 0, 'wght' 400, 'GRAD' 0, 'opsz' 24; }
mark { background: var(--mark); color: inherit; padding: 0 1px; }
.layout { min-height: 100vh; display: flex; flex-direction: column; }
.topbar { flex: none; position: sticky; top: 0; z-index: 130; height: var(--topbar-h); display: flex; align-items: center; justify-content: space-between; gap: 12px; padding: 0 16px; background: var(--panel); border-bottom: 1px solid var(--border); }
.brand { display: flex; align-items: center; gap: 10px; font-weight: 700; white-space: nowrap; color: inherit; text-decoration: none; }
.logo { width: 26px; height: 26px; flex: none; display: block; }
.body-row { flex: 1; display: flex; }
.sidebar { flex: none; width: 224px; background: var(--panel); border-right: 1px solid var(--border); padding: 14px 8px 16px; }
.nav-section-head { font-size: 11px; font-weight: 700; letter-spacing: .06em; text-transform: uppercase; color: var(--muted); padding: 0 8px 6px; }
.nav-item { display: flex; align-items: center; gap: 9px; padding: 7px 8px; font-size: 14px; color: var(--text); text-decoration: none; }
.nav-item.active { background: color-mix(in srgb, var(--accent) 14%, transparent); color: var(--accent); }
.nav-item .material-symbols-rounded { font-size: 19px; color: var(--muted); }
.nav-item.active .material-symbols-rounded { color: var(--accent); }
.content { flex: 1; padding: 18px 20px; max-width: 1020px; }
/* query bar */
.query-bar { display: flex; gap: 10px; align-items: stretch; margin-bottom: 6px; position: relative; }
.query-input { flex: 1; display: flex; align-items: center; gap: 10px; background: var(--panel); border: 1px solid var(--border); padding: 0 12px; height: 40px; }
.query-input:focus-within { border-color: var(--accent); box-shadow: 0 0 0 1px var(--accent); }
.query-input .material-symbols-rounded { color: var(--muted); font-size: 20px; }
.query-input input { flex: 1; background: none; border: none; outline: none; color: var(--text); font-size: 15px; height: 100%; }
.seg { display: flex; border: 1px solid var(--border); background: var(--panel); }
.seg button { appearance: none; border: none; background: none; color: var(--muted); font-size: 13px; font-weight: 600; padding: 0 14px; cursor: pointer; border-right: 1px solid var(--border); }
.seg button:last-child { border-right: none; }
.seg button.on { background: color-mix(in srgb, var(--accent) 16%, transparent); color: var(--accent); box-shadow: inset 0 -2px 0 var(--accent); }
/* roots multiselect */
.ms { position: relative; }
.ms-btn { display: flex; align-items: center; gap: 8px; background: var(--panel); border: 1px solid var(--accent); color: var(--text); font-size: 13px; height: 40px; padding: 0 12px; cursor: pointer; min-width: 130px; justify-content: space-between; }
.ms-btn .material-symbols-rounded { font-size: 17px; color: var(--muted); }
.ms-drop { position: absolute; top: 44px; left: 0; width: 230px; background: var(--panel); border: 1px solid var(--border); box-shadow: 0 12px 32px rgba(0,0,0,.5); z-index: 40; padding: 6px 0; }
.ms-opt { display: flex; align-items: center; gap: 10px; padding: 7px 12px; font-size: 14px; cursor: pointer; }
.ms-opt:hover { background: color-mix(in srgb, var(--panel) 50%, var(--bg)); }
.ms-opt .box { width: 16px; height: 16px; border: 1px solid var(--muted); display: inline-flex; align-items: center; justify-content: center; flex: none; }
.ms-opt .box.on { background: var(--accent); border-color: var(--accent); }
.ms-opt .box .material-symbols-rounded { font-size: 13px; color: #fff; }
.ms-opt .rp { margin-left: auto; color: var(--muted); font-size: 12px; }
/* stop button (visible while searching) */
.btn-stop { display: inline-flex; align-items: center; gap: 7px; border: 1px solid var(--danger-border); color: var(--danger); background: color-mix(in srgb, var(--danger) 8%, var(--panel)); font-size: 13px; font-weight: 600; padding: 0 14px; height: 40px; cursor: pointer; }
.btn-stop .material-symbols-rounded { font-size: 17px; font-variation-settings: 'FILL' 1; }
.btn { display: inline-flex; align-items: center; gap: 7px; border: 1px solid var(--border); background: var(--panel); color: var(--text); font-size: 13px; font-weight: 600; padding: 0 14px; height: 40px; cursor: pointer; }
.btn .material-symbols-rounded { font-size: 17px; }
/* live status line */
.status-line { display: flex; align-items: center; gap: 8px; margin: 4px 2px 12px; font-size: 12.5px; color: var(--muted); }
.status-line .dot { width: 8px; height: 8px; background: var(--accent); animation: pulse 1.2s ease-in-out infinite; }
@keyframes pulse { 50% { opacity: .25; } }
.status-line .grow { flex: 1; }
.status-line .linkish { color: var(--muted); cursor: pointer; }
/* result groups */
.group-label { display: flex; align-items: center; gap: 7px; margin: 16px 2px 8px; font-size: 12px; font-weight: 700; letter-spacing: .05em; text-transform: uppercase; color: var(--muted); }
.group-label .material-symbols-rounded { font-size: 15px; }
.group-label .n { font-weight: 400; letter-spacing: 0; text-transform: none; }
.group-label .grow { flex: 1; }
.collapse-all { appearance: none; border: none; background: none; color: var(--muted); font-size: 12px; display: inline-flex; align-items: center; gap: 4px; cursor: pointer; padding: 2px 4px; }
.collapse-all:hover { color: var(--text); }
.collapse-all .material-symbols-rounded { font-size: 15px; }
.entries-list { background: var(--panel); border: 1px solid var(--border); overflow: hidden; }
.srow { display: grid; grid-template-columns: 28px minmax(0,1fr) 240px 90px; gap: 12px; align-items: center; padding: 9px 14px; border-bottom: 1px solid var(--border); font-size: 14px; cursor: pointer; }
.srow:last-child { border-bottom: none; }
.srow:hover { background: color-mix(in srgb, var(--panel) 50%, var(--bg)); }
.srow .material-symbols-rounded { font-size: 20px; color: var(--muted); }
.spath { color: var(--muted); font-size: 13px; overflow: hidden; text-overflow: ellipsis; white-space: nowrap; }
.ssize { color: var(--muted); font-size: 13px; text-align: right; }
.more-row { display: flex; justify-content: center; gap: 8px; padding: 7px; color: var(--muted); font-size: 13px; border-top: 1px solid var(--border); }
/* content match cards — match lines only */
.match-card { background: var(--panel); border: 1px solid var(--border); margin-bottom: 10px; }
.match-head { display: grid; grid-template-columns: 28px minmax(0,1fr) 70px 90px 34px; gap: 12px; align-items: center; padding: 9px 10px 9px 14px; font-size: 14px; }
.match-head .material-symbols-rounded.mi { font-size: 20px; color: var(--muted); }
.match-head b { font-weight: 600; }
.match-head .mpath { overflow: hidden; text-overflow: ellipsis; white-space: nowrap; }
.match-head .mcount { color: var(--accent); font-size: 12px; font-weight: 600; text-align: right; }
.match-head .msize { color: var(--muted); font-size: 13px; text-align: right; }
.chev { appearance: none; border: 1px solid transparent; background: none; color: var(--muted); cursor: pointer; padding: 5px; display: inline-flex; align-items: center; justify-content: center; }
.chev:hover { border-color: var(--border); background: var(--bg); color: var(--text); }
.chev .material-symbols-rounded { font-size: 18px; }
.match-lines { border-top: 1px solid var(--border); background: color-mix(in srgb, var(--bg) 60%, var(--panel)); padding: 6px 0; }
.mline { display: grid; grid-template-columns: 52px minmax(0,1fr); gap: 12px; padding: 3px 14px 3px 10px; font-family: "JetBrains Mono", monospace; font-size: 12.5px; line-height: 1.5; cursor: pointer; }
.mline:hover { background: color-mix(in srgb, var(--accent) 10%, transparent); }
.mline .ln { color: var(--muted); text-align: right; user-select: none; }
.mline .tx { white-space: pre; overflow: hidden; text-overflow: ellipsis; }
</style>
</head>
<body>
<div class="layout">
<header class="topbar">
<a class="brand" href="#">
<svg class="logo" viewBox="0 0 32 32"><rect width="32" height="32" fill="#2563eb"/><path d="M6 8h7l3 4h10v14H6z" fill="#fff"/><path d="M11.5 17.5 16 22l4.5-4.5" fill="none" stroke="#2563eb" stroke-width="3"/></svg>
<span>filebrowser-ng</span>
</a>
</header>
<div class="body-row">
<aside class="sidebar">
<div class="nav-section-head">Files</div>
<div class="nav-item active"><span class="material-symbols-rounded">search</span> Search</div>
<div class="nav-item"><span class="material-symbols-rounded">folder</span> Documents</div>
<div class="nav-item"><span class="material-symbols-rounded">folder</span> Projects</div>
<div class="nav-section-head" style="margin-top:14px">Manage</div>
<div class="nav-item"><span class="material-symbols-rounded">link</span> Shares</div>
<div class="nav-item"><span class="material-symbols-rounded">settings</span> Settings</div>
</aside>
<main class="content">
<div class="query-bar">
<div class="query-input">
<span class="material-symbols-rounded">search</span>
<input value="invoice 2025" spellcheck="false">
</div>
<div class="seg">
<button>Name</button>
<button>Content</button>
<button class="on">Both</button>
</div>
<div class="ms">
<button class="ms-btn"><span>2 roots</span><span class="material-symbols-rounded">expand_more</span></button>
<div class="ms-drop">
<div class="ms-opt"><span class="box on"><span class="material-symbols-rounded">check</span></span> Documents <span class="rp">612 files</span></div>
<div class="ms-opt"><span class="box on"><span class="material-symbols-rounded">check</span></span> Projects <span class="rp">588 files</span></div>
</div>
</div>
<button class="btn-stop"><span class="material-symbols-rounded">stop</span> Stop</button>
</div>
<div class="status-line">
<span class="dot"></span>
<span>searching… scanned 412 files, 5 name matches, 14 content matches</span>
<span class="grow"></span>
</div>
<div class="group-label"><span class="material-symbols-rounded">description</span> Files <span class="n">· 5</span></div>
<div class="entries-list">
<div class="srow"><span class="material-symbols-rounded">description</span><span><mark>invoice</mark>_2025-11_final.pdf</span><span class="spath">Documents / finance / 2025</span><span class="ssize">182 KB</span></div>
<div class="srow"><span class="material-symbols-rounded">description</span><span><mark>invoice</mark>_2025-08_draeger.pdf</span><span class="spath">Documents / finance / 2025</span><span class="ssize">164 KB</span></div>
<div class="srow"><span class="material-symbols-rounded">folder</span><span><mark>invoices</mark>-2025</span><span class="spath">Projects / shop</span><span class="ssize">—</span></div>
<div class="srow"><span class="material-symbols-rounded">grid_on</span><span>export_2025_Q3.<mark>xlsx</mark></span><span class="spath">Documents / finance</span><span class="ssize">58 KB</span></div>
</div>
<div class="more-row"><span class="material-symbols-rounded" style="font-size:16px">more_horiz</span> 1 more…</div>
<div class="group-label">
<span class="material-symbols-rounded">text_snippet</span> Content <span class="n">· 14 matches in 4 files</span>
<span class="grow"></span>
<button class="collapse-all"><span class="material-symbols-rounded">unfold_less</span> Collapse all</button>
</div>
<div class="match-card">
<div class="match-head">
<span class="material-symbols-rounded mi">code</span>
<span class="mpath">Projects / billing / <b>generate_2025.rs</b></span>
<span class="msize">14 KB</span>
<span class="mcount">6 matches</span>
<button class="chev" title="Collapse"><span class="material-symbols-rounded">expand_more</span></button>
</div>
<div class="match-lines">
<div class="mline"><span class="ln">42</span><span class="tx">let template = load("<mark>invoice</mark>-<mark>2025</mark>-layout");</span></div>
<div class="mline"><span class="ln">58</span><span class="tx">stamp(&mut pdf, "<mark>INVOICE</mark>-<mark>2025</mark>-…");</span></div>
<div class="mline"><span class="ln">87</span><span class="tx">// <mark>2025</mark>-11: VAT rate changed mid-run</span></div>
<div class="mline"><span class="ln">91</span><span class="tx">fn total_2025(rows: &[Line]) -> Money { … }</span></div>
</div>
</div>
<div class="match-card">
<div class="match-head">
<span class="material-symbols-rounded mi">description</span>
<span class="mpath">Documents / finance / <b>2025_summary.md</b></span>
<span class="msize">6 KB</span>
<span class="mcount">5 matches</span>
<button class="chev" title="Collapse"><span class="material-symbols-rounded">expand_more</span></button>
</div>
<div class="match-lines">
<div class="mline"><span class="ln">3</span><span class="tx">## <mark>Invoice</mark> totals <mark>2025</mark></span></div>
<div class="mline"><span class="ln">11</span><span class="tx">| <mark>2025</mark>-11 | 4,120.00 | paid |</span></div>
<div class="mline"><span class="ln">12</span><span class="tx">| <mark>2025</mark>-12 | 3,807.50 | open |</span></div>
</div>
</div>
<div class="match-card">
<div class="match-head">
<span class="material-symbols-rounded mi">grid_on</span>
<span class="mpath">Documents / finance / <b>export_2025_Q3.csv</b></span>
<span class="msize">210 KB</span>
<span class="mcount">3 matches</span>
<button class="chev" title="Expand"><span class="material-symbols-rounded">expand_less</span></button>
</div>
</div>
</main>
</div>
</div>
</body>
</html>
Mapi-types/src/lib.rs
@@ -50,8 +50,8 @@ pub const P_SHARE: &str = "share";
pub const P_Q: &str = "q";
/// Which index to search: `name`, `content` or `both`.
pub const P_SCOPE: &str = "scope";
/// Comma-separated root ids; omitted = all of the user's roots.
pub const P_ROOTS: &str = "roots";
/// Root id to search; omitted = the caller's first root.
pub const P_ROOT: &str = "root";
/// `?overwrite=true|1` on mutations and uploads.
pub const P_OVERWRITE: &str = "overwrite";
@@ -237,6 +237,9 @@ pub enum SearchEvent {
path: String,
size: u64,
is_dir: bool,
/// Sniffed the same way as a directory listing's, so the client can
/// pick an icon and a viewer without a second guess at the name.
kind: FileKind,
},
/// One matching line (scope `content`/`both`). `path` is relative to the
/// root; `text` is the matched line, truncated to a fixed length.
Mserver/src/api/search.rs
@@ -1,6 +1,6 @@
//! `GET /api/search` — index-free name and content search, streamed as SSE.
//!
//! No index, by design: every search walks the selected roots with
//! No index, by design: every search walks the selected root with
//! [`ignore::WalkParallel`] (fd/ripgrep's walker: parallel, skips hidden
//! files, honors `.gitignore`) and matches on the fly.
//!
@@ -14,6 +14,9 @@
//! 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
@@ -28,16 +31,15 @@ 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 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, HasState};
use crate::api::common::AuthUser;
use crate::db::RootRow;
use crate::error::{ApiError, AppState};
@@ -59,40 +61,25 @@ pub(super) struct SearchQuery {
/// `name` (default), `content` or `both`.
#[serde(default)]
scope: Option<String>,
/// Comma-separated root ids; empty = all of the user's roots.
/// The root to search; omitted = the caller's first root.
#[serde(default)]
roots: Option<String>,
}
/// 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<S> axum::extract::FromRequestParts<S> for SearchUser
where
S: HasState + Send + Sync,
{
type Rejection = ApiError;
async fn from_request_parts(parts: &mut Parts, state: &S) -> Result<Self, Self::Rejection> {
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))
}
root: Option<i64>,
}
pub(super) async fn search(
State(state): State<Arc<AppState>>,
SearchUser(auth): SearchUser,
auth: AuthUser,
AxumQuery(query): AxumQuery<SearchQuery>,
) -> Result<Response, ApiError> {
// 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()
@@ -117,84 +104,53 @@ pub(super) async fn search(
));
}
// 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<RootRow, ApiError> {
let id = s.parse::<i64>().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::<Result<Vec<RootRow>, ApiError>>(),
}?;
Ok(sse_response(search_stream(
// 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",
)
})?;
let events = search_stream(
q.to_string(),
want_name,
want_content,
Arc::new(roots),
root,
state.root.clone(),
)))
}
/// Wraps an event stream in an SSE response (`data: <json>` per event).
fn sse_response(
events: impl futures_util::Stream<Item = SearchEvent> + Send + 'static,
) -> Response {
let body = events.map(|ev| {
// Infallible: SearchEvent is always serializable.
let json = serde_json::to_string(&ev).expect("SearchEvent serializes");
Ok::<String, std::io::Error>(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")
);
// `Sse` does the `data: <json>\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. The counters are `Arc`'d so the per-thread
/// walk visitors (which must be `Clone`) can update them lock-free.
/// 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<SearchEvent>,
words: Vec<String>,
roots: Arc<Vec<RootRow>>,
root: RootRow,
server_root: PathBuf,
started: Instant,
/// Directory entries visited (both phases).
scanned: Arc<AtomicUsize>,
/// Files skipped by the content phase (over the size cap).
skipped: Arc<AtomicUsize>,
names: Arc<AtomicUsize>,
matches: Arc<AtomicUsize>,
/// Directory entries visited.
scanned: AtomicUsize,
/// Files skipped for content search (over the size cap).
skipped: AtomicUsize,
names: AtomicUsize,
matches: AtomicUsize,
}
/// 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
/// 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
@@ -204,31 +160,24 @@ fn search_stream(
q: String,
want_name: bool,
want_content: bool,
roots: Arc<Vec<RootRow>>,
root: RootRow,
server_root: PathBuf,
) -> impl futures_util::Stream<Item = SearchEvent> + Send {
let (tx, rx) = tokio::sync::mpsc::channel::<SearchEvent>(256);
let state = Arc::new(SearchState {
tx,
words: q.split_whitespace().map(|w| w.to_lowercase()).collect(),
roots,
root,
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)),
scanned: AtomicUsize::new(0),
skipped: AtomicUsize::new(0),
names: AtomicUsize::new(0),
matches: 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);
}
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()
@@ -246,131 +195,110 @@ fn search_stream(
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<F>(root: &RootRow, server_root: &Path, visit: F)
where
F: FnMut(&DirEntry, &Path) -> WalkState + Clone + Send,
{
let abs = server_root.join(&root.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<SearchState>, q: &str, want_name: bool, want_content: bool) {
let abs = st.server_root.join(&st.root.path);
if !abs.is_dir() {
// Root removed out from under us: skip, keep searching the rest.
// Root removed out from under us: nothing to search.
return;
}
WalkBuilder::new(&abs)
.standard_filters(true) // fd's defaults: hidden files + gitignore
.build_parallel()
.run(move || {
let mut visit = visit.clone();
.run(|| {
let abs = abs.clone();
Box::new(move |result| match result {
Ok(entry) => visit(&entry, &abs),
Ok(entry) => visit(st, q, want_name, want_content, &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<SearchState>) {
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('\\', "/");
/// One walked entry: emit a name hit, grep it, or both.
fn visit(
st: &Arc<SearchState>,
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 root itself has an empty relative path and is not a result.
let Some(rel) = entry
.path()
.strip_prefix(abs)
.ok()
.filter(|p| !p.as_os_str().is_empty())
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());
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 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,
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.
match st.tx.blocking_send(ev) {
Ok(()) => WalkState::Continue,
Err(_) => WalkState::Quit,
if st.tx.blocking_send(ev).is_err() {
return WalkState::Quit;
}
});
}
}
}
/// Phase 2: content matches.
fn content_phase(st: &Arc<SearchState>, 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
}
});
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: &Path,
root_id: i64,
tx: &tokio::sync::mpsc::Sender<SearchEvent>,
matches: &Arc<AtomicUsize>,
) {
fn search_one_file(q: &str, path: &Path, rel: String, st: &SearchState) {
let matcher = match RegexMatcherBuilder::new()
.case_insensitive(true)
.build_literals(&[q])
@@ -382,14 +310,7 @@ fn search_one_file(
// 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,
};
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);
}
@@ -398,11 +319,9 @@ fn search_one_file(
/// 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<SearchEvent>,
matches: &'a Arc<AtomicUsize>,
st: &'a SearchState,
/// Matched lines emitted for the current file.
lines: usize,
root_id: i64,
rel: String,
}
@@ -410,7 +329,7 @@ impl Sink for MatchSink<'_> {
type Error = std::io::Error;
fn matched(&mut self, _searcher: &Searcher, mat: &SinkMatch) -> Result<bool, Self::Error> {
if self.tx.is_closed() || self.lines >= SEARCH_MAX_LINES_PER_FILE {
if self.st.tx.is_closed() || self.lines >= SEARCH_MAX_LINES_PER_FILE {
return Ok(false);
}
let line = match mat.line_number() {
@@ -424,13 +343,13 @@ impl Sink for MatchSink<'_> {
};
// 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 {
if self.st.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,
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,
@@ -443,23 +362,17 @@ impl Sink for MatchSink<'_> {
}
}
/// Strips the line terminator and truncates at [`SEARCH_MAX_LINE_BYTES`],
/// backing off to a UTF-8 boundary if the cut split a character.
/// 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(|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()
.map_or(0, |i| i + 1);
String::from_utf8_lossy(&line[..end.min(SEARCH_MAX_LINE_BYTES)]).into_owned()
}
#[cfg(test)]
Aserver/tests/api_search.rs
@@ -0,0 +1,96 @@
//! Search API: one walk serves both scopes, streamed as SSE.
mod common;
use axum::http::StatusCode;
use common::*;
/// Root id for the whole-root (".") user root is 1 (first row inserted).
const ROOT: i64 = 1;
/// The `data:` payloads of an SSE body, in order.
fn events(body: &str) -> Vec<serde_json::Value> {
body.lines()
.filter_map(|l| l.strip_prefix("data:"))
.map(|d| serde_json::from_str(d.trim()).expect("event is JSON"))
.collect()
}
#[tokio::test]
async fn both_scopes_stream_from_one_walk() {
let env = Env::new().await;
let admin = env.admin().await;
let r = admin
.get(&format!("/api/search?q=hello&scope=both&root={ROOT}"))
.await;
assert_eq!(r.status, StatusCode::OK);
assert_eq!(
r.header("content-type").as_deref(),
Some("text/event-stream")
);
assert_eq!(r.header("x-accel-buffering").as_deref(), Some("no"));
let evs = events(&r.text());
// The name hit: matched on the entry's own name, with the server's
// sniffed kind.
let file = evs
.iter()
.find(|e| e["type"] == "file")
.expect("a name hit");
assert_eq!(file["path"], "docs/inner/hello.txt");
assert_eq!(file["kind"], "text");
assert_eq!(file["is_dir"], false);
// The content hit, from the same walk.
let m = evs
.iter()
.find(|e| e["type"] == "match")
.expect("a content hit");
assert_eq!(m["path"], "docs/inner/hello.txt");
assert_eq!(m["line"], 1);
assert_eq!(m["text"], "hello world");
// The stream always ends with the summary.
let done = evs.last().expect("a done event");
assert_eq!(done["type"], "done");
assert_eq!(done["stopped"], false);
assert!(done["scanned"].as_u64().unwrap() >= 8);
}
#[tokio::test]
async fn name_scope_ignores_the_parent_path() {
let env = Env::new().await;
let admin = env.admin().await;
let r = admin
.get(&format!("/api/search?q=docs&scope=name&root={ROOT}"))
.await;
let paths: Vec<String> = events(&r.text())
.iter()
.filter(|e| e["type"] == "file")
.map(|e| e["path"].as_str().unwrap().to_string())
.collect();
// The folder itself, never the files inside it.
assert_eq!(paths, vec!["docs".to_string()]);
}
#[tokio::test]
async fn search_rejects_bad_input_and_anonymous_callers() {
let env = Env::new().await;
let admin = env.admin().await;
let anon = Client::new(env.app.clone());
assert_eq!(
anon.get("/api/search?q=hello").await.status,
StatusCode::UNAUTHORIZED
);
assert_eq!(
admin.get("/api/search?q=%20").await.status,
StatusCode::BAD_REQUEST
);
assert_eq!(
admin.get("/api/search?q=hello&scope=nope").await.status,
StatusCode::BAD_REQUEST
);
assert_eq!(
admin.get("/api/search?q=hello&root=999").await.status,
StatusCode::FORBIDDEN
);
}
Mweb/Cargo.toml
@@ -15,7 +15,6 @@ thiserror = "2"
wasm-bindgen = "0.2"
wasm-bindgen-futures = "0.4"
web-sys = { version = "0.3", features = [
"AbortController",
"AddEventListenerOptions",
"Blob",
"Clipboard",
@@ -23,6 +22,7 @@ web-sys = { version = "0.3", features = [
"DomRect",
"Element",
"Event",
"EventSource",
"EventTarget",
"EventListenerOptions",
"File",
@@ -37,12 +37,9 @@ web-sys = { version = "0.3", features = [
"HtmlSelectElement",
"KeyboardEvent",
"Location",
"MessageEvent",
"MouseEvent",
"Navigator",
"Performance",
"ReadableStream",
"ReadableStreamDefaultReader",
"ReadableStreamReadResult",
"Request",
"RequestInit",
"RequestMode",
@@ -50,7 +47,6 @@ web-sys = { version = "0.3", features = [
"Storage",
"SubmitEvent",
"SvgElement",
"TextDecodeOptions",
"TextDecoder",
"UrlSearchParams",
"Window",
] }
Mweb/src/api.rs
@@ -9,12 +9,13 @@ use serde::Serialize;
use serde::de::DeserializeOwned;
use wasm_bindgen::JsCast;
use wasm_bindgen::JsValue;
use wasm_bindgen::closure::Closure;
use wasm_bindgen_futures::JsFuture;
use api_types::{
ACTION_CONTENT, ACTION_CREATE_FILE, ACTION_DOWNLOAD, ACTION_MKDIR, ACTION_PREVIEW,
ADMIN_SETTINGS, ADMIN_USERS, AUTH_LOGIN, AUTH_LOGOUT, AUTH_ME, AUTH_SETUP, CreateShare,
CreateUser, Credentials, FILES, Mutation, P_ACTION, P_FORMAT, P_OVERWRITE, P_Q, P_ROOTS,
CreateUser, Credentials, FILES, Mutation, P_ACTION, P_FORMAT, P_OVERWRITE, P_Q, P_ROOT,
P_SCOPE, P_SHARE, Root, SEARCH, SHARE, SHARES, Settings, UpdateUser,
};
pub use api_types::{
@@ -587,8 +588,6 @@ fn random_boundary_suffix() -> String {
s
}
use wasm_bindgen::closure::Closure;
/// Open the native file dialog and run `on_files` with the picked files once
/// the user confirms. `directory` uses webkitdirectory (folder upload).
fn webkit_relative_path(file: &web_sys::File) -> String {
@@ -759,127 +758,54 @@ async fn parse_error_body(resp: &web_sys::Response) -> ErrBody {
/// Start a search and stream its results as they are found.
///
/// The server answers with an SSE stream (one `data: <json>` event per
/// result, ending in a `done` event). `on_event` runs once per event on the
/// main thread; the returned [`AbortController`] stops the search client-side
/// (the server notices the dropped connection and unwinds its walk).
/// [`web_sys::EventSource`] is the platform's SSE client: it frames the
/// stream and parses the events. `on_event` runs once per event on the main
/// thread; `.close()` on the returned source stops the search (the server
/// notices the dropped connection and unwinds its walk).
///
/// `roots` empty = all of the caller's roots. Returns an error before
/// streaming starts only for request/HTTP failures; a stream that dies
/// mid-way simply ends without a `done` event.
/// A stream that ends or dies mid-way simply stops, without a `done` event.
/// An `EventSource` cannot read the server's JSON error body, so a rejected
/// request reports one generic message.
pub fn search_stream(
q: String,
scope: &str,
roots: &[i64],
root: i64,
on_event: leptos::prelude::Callback<api_types::SearchEvent, ()>,
on_error: leptos::prelude::Callback<String, ()>,
) -> Result<web_sys::AbortController, ApiError> {
let mut url = format!(
"{}?{P_Q}={}&{P_SCOPE}={}",
SEARCH,
) -> Result<web_sys::EventSource, ApiError> {
let url = format!(
"{SEARCH}?{P_Q}={}&{P_SCOPE}={scope}&{P_ROOT}={root}",
js_sys::encode_uri_component(&q),
scope
);
if !roots.is_empty() {
url.push_str(&format!(
"&{P_ROOTS}={}",
roots
.iter()
.map(|r| r.to_string())
.collect::<Vec<_>>()
.join(",")
));
}
let src = web_sys::EventSource::new(&url).map_err(|e| ApiError::Net(format!("{e:?}")))?;
let controller = web_sys::AbortController::new()
.map_err(|_| ApiError::Net("AbortController unavailable".to_string()))?;
let opts = web_sys::RequestInit::new();
opts.set_mode(web_sys::RequestMode::SameOrigin);
opts.set_signal(Some(&controller.signal()));
wasm_bindgen_futures::spawn_local(async move {
// HTTP-level failure (4xx/5xx JSON error from the server).
let resp = match fetch_checked(&url, &opts, "search failed").await {
Ok(r) => r,
Err(e) => {
// Aborting (user pressed Stop) rejects the fetch with an
// AbortError; that is expected, not a failure.
let msg = e.to_string();
if !msg.contains("AbortError") {
on_error.run(msg);
}
let on_msg =
Closure::<dyn FnMut(web_sys::MessageEvent)>::new(move |ev: web_sys::MessageEvent| {
let Some(data) = ev.data().as_string() else {
return;
}
};
let body = match resp.body() {
Some(b) => b,
None => return,
};
let reader = match body
.get_reader()
.dyn_into::<web_sys::ReadableStreamDefaultReader>()
{
Ok(r) => r,
Err(_) => return,
};
let decoder = match web_sys::TextDecoder::new() {
Ok(d) => d,
Err(_) => return,
};
let decode_opts = web_sys::TextDecodeOptions::new();
// Streaming mode: a multi-byte character split across two chunks
// must not be mangled.
decode_opts.set_stream(true);
let mut buf = String::new();
loop {
let result = match JsFuture::from(reader.read()).await {
Ok(v) => v,
Err(_) => return, // aborted or network failure: stop quietly
};
// The read() result is a plain { done, value } dictionary, so
// read its fields instead of downcasting to the (branded)
// ReadableStreamReadResult type.
let done =
js_sys::Reflect::get(&result, &JsValue::from_str("done")).unwrap_or(JsValue::TRUE);
if done.as_bool().unwrap_or(true) {
break;
}
let value = match js_sys::Reflect::get(&result, &JsValue::from_str("value")) {
Ok(v) => v,
Err(_) => break,
};
let u8: js_sys::Uint8Array = match value.dyn_into() {
Ok(u8) => u8,
Err(_) => break,
};
let chunk = match decoder.decode_with_u8_array(u8.to_vec().as_slice()) {
Ok(c) => c,
Err(_) => break,
};
buf.push_str(&chunk);
// SSE framing: events are separated by a blank line.
while let Some(pos) = buf.find("\n\n") {
let event = buf[..pos].to_string();
buf.drain(..pos + 2);
handle_sse_event(&event, &on_event);
if let Ok(ev) = serde_json::from_str::<api_types::SearchEvent>(&data) {
on_event.run(ev);
}
});
src.set_onmessage(Some(on_msg.as_ref().unchecked_ref()));
on_msg.forget();
// `error` fires both when the request is rejected (`CLOSED`) and when the
// stream ends normally (`CONNECTING`, a reconnect pending). A reconnect
// would re-run the whole search, so close the source either way.
let s = src.clone();
let on_err = Closure::<dyn FnMut(web_sys::Event)>::new(move |_| {
let failed = s.ready_state() == web_sys::EventSource::CLOSED;
s.close();
if failed {
on_error.run(crate::i18n::error_text(None, "search failed"));
}
});
src.set_onerror(Some(on_err.as_ref().unchecked_ref()));
on_err.forget();
Ok(controller)
}
/// Parse one SSE block (`data: {json}`) and forward the event.
fn handle_sse_event(event: &str, on_event: &leptos::prelude::Callback<api_types::SearchEvent, ()>) {
for line in event.lines() {
let Some(data) = line.strip_prefix("data:") else {
continue;
};
if let Ok(ev) = serde_json::from_str::<api_types::SearchEvent>(data.trim()) {
on_event.run(ev);
}
}
Ok(src)
}
/// Read a form input's value by element id.
Mweb/src/i18n.rs
@@ -456,7 +456,6 @@ pub mod k {
pub const ERR_FS_CONFLICT: &str = "err_fs_conflict";
pub const ERR_SEARCH_FORBIDDEN: &str = "err_search_forbidden";
pub const ERR_SEARCH_BAD_SCOPE: &str = "err_search_bad_scope";
pub const ERR_SEARCH_BAD_ROOTS: &str = "err_search_bad_roots";
pub const ERR_ROOT_FORBIDDEN: &str = "err_root_forbidden";
}
@@ -834,7 +833,6 @@ const EN: &[(&str, &str)] = &[
k::ERR_SEARCH_BAD_SCOPE,
"scope must be name, content or both",
),
(k::ERR_SEARCH_BAD_ROOTS, "root ids must be numbers"),
(k::ERR_ROOT_FORBIDDEN, "root not accessible"),
];
@@ -1271,7 +1269,6 @@ const DE: &[(&str, &str)] = &[
k::ERR_SEARCH_BAD_SCOPE,
"Bereich muss name, content oder both sein",
),
(k::ERR_SEARCH_BAD_ROOTS, "Ordner-IDs müssen Zahlen sein"),
(k::ERR_ROOT_FORBIDDEN, "Ordner nicht verfügbar"),
];
@@ -1714,10 +1711,6 @@ const FR: &[(&str, &str)] = &[
k::ERR_SEARCH_BAD_SCOPE,
"le périmètre doit être name, content ou both",
),
(
k::ERR_SEARCH_BAD_ROOTS,
"les identifiants de dossiers doivent être des nombres",
),
(k::ERR_ROOT_FORBIDDEN, "dossier inaccessible"),
];
Mweb/src/views/search.rs
@@ -1,13 +1,13 @@
//! The search view (`#/search`): index-free name and content search with
//! streamed results, an explicit stop button, a roots multi-select, and
//! collapsible per-file match lists.
//! streamed results, an explicit stop button, a root select, and collapsible
//! per-file match lists.
//!
//! Results arrive as they are found (SSE, see `api::search_stream`). The
//! client caps only what it *renders* (name rows a page of [`FILE_PAGE`] at a
//! time, match lines per file at [`RENDER_LINE_CAP`]) — the search itself is
//! unbounded and ends when the walk is done or the user stops it. Stopping
//! aborts the fetch; the server notices the dropped connection and unwinds
//! its walk.
//! closes the event source; the server notices the dropped connection and
//! unwinds its walk.
use std::collections::HashMap;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
@@ -53,47 +53,17 @@ const RETAIN_MATCH_FILE_CAP: usize = 500;
/// per event: an open search can emit thousands of events and applying them
/// one by one would re-render per event and peg the main thread.
///
/// A flush applies its whole batch, but in feedback-sized chunks that yield
/// to the browser between them (see [`apply_chunked`]) rather than in one go.
/// A flush applies its whole batch, but in chunks that yield to the browser
/// between them (see [`apply_chunked`]) rather than in one go.
const FLUSH_MS: u32 = 120;
/// Target and back-off thresholds for one chunk plus its yield.
///
/// Chunks are sized by feedback, not by a fixed item count: the cost of a row
/// depends on the machine, and a hardcoded count either blocks on a slow one
/// or trickles for seconds on a fast one.
/// Events folded between yields, and items written per chunk.
///
/// The signal is the interval between successive resumptions of the render
/// loop, which covers the chunk's own work, the `<For>` pass it triggers, and
/// anything else the browser did in between. Keeping that near
/// [`CHUNK_TARGET_MS`] keeps the main thread free most of the time, so input
/// stays responsive while a large result set fills in.
const CHUNK_TARGET_MS: f64 = 10.0;
const CHUNK_SLOW_MS: f64 = 14.0;
/// Events folded between yields. Fixed rather than adaptive: folding an
/// event is a bounded allocate-and-hash, so unlike rendering it has no
/// machine-dependent cliff worth measuring for.
// ponytail: two fixed sizes, not a feedback controller. `<For>` re-runs
// `each` over the whole list per write, so a write costs roughly
// `CHUNK + list_len` and a slow machine can still overrun a frame. Measure
// before replacing this with anything adaptive.
const EVENTS_PER_SLICE: usize = 2048;
/// First chunk's size, before there is a measurement to go on.
const INITIAL_CHUNK: usize = 64;
/// Chunk bounds. The floor keeps forward progress when every item is
/// expensive; the ceiling stops one chunk from blocking when the interval
/// looks cheap for an unrelated reason.
/// The ceiling matters more than it looks: `<For>` re-runs `each` over the
/// whole list per write, so a write costs roughly `chunk + list_len`. As the
/// list fills, a fixed chunk gets steadily more expensive, and the feedback
/// below can only react *after* a round overran. Capping the chunk caps how
/// bad that one round can be.
const MIN_CHUNK: usize = 16;
const MAX_CHUNK: usize = 128;
/// `performance.now()`, or 0 when there is no window (never in the browser).
fn now_ms() -> f64 {
web_sys::window()
.and_then(|w| w.performance())
.map(|p| p.now())
.unwrap_or(0.0)
}
const CHUNK: usize = 64;
/// Hands the main thread back to the browser, then resumes.
///
@@ -115,24 +85,6 @@ async fn yield_to_browser() {
TimeoutFuture::new(0).await;
}
/// Next chunk size from the last round's interval.
///
/// Additive increase, multiplicative decrease. Doubling on the way up
/// overshoots: it keeps growing until a chunk finally overruns, and *that*
/// chunk is the visible stall (measured as a 241 ms spike once the size ran
/// away to 1024). Growing an eighth at a time bounds the overshoot, while
/// halving still backs off from a bad round immediately.
fn next_chunk(prev: usize, interval_ms: f64) -> usize {
let next = if interval_ms > CHUNK_SLOW_MS {
prev / 2
} else if interval_ms < CHUNK_TARGET_MS {
prev + (prev / 8).max(8)
} else {
prev
};
next.clamp(MIN_CHUNK, MAX_CHUNK)
}
/// Applies `items` to `write` in chunks, yielding to the browser between
/// them. One signal write per chunk, so a chunk costs one `<For>` pass.
async fn apply_chunked<T, F>(items: Vec<T>, my_gen: u64, ctx: &FlushCtx, mut write: F)
@@ -140,20 +92,14 @@ where
F: FnMut(Vec<T>),
{
let mut rest = items;
let mut chunk = INITIAL_CHUNK;
let mut last_round = now_ms();
while !rest.is_empty() {
// A new search (or an unmount) makes the remaining items obsolete.
if ctx.guard.load(Ordering::Relaxed) || ctx.search_gen.load(Ordering::SeqCst) != my_gen {
if ctx.stale(my_gen) {
return;
}
let take = chunk.min(rest.len());
let head: Vec<T> = rest.drain(..take).collect();
write(head);
let take = CHUNK.min(rest.len());
write(rest.drain(..take).collect());
yield_to_browser().await;
let now = now_ms();
chunk = next_chunk(chunk, now - last_round);
last_round = now;
}
}
@@ -213,35 +159,28 @@ struct FileHit {
path: Arc<str>,
size: u64,
is_dir: bool,
kind: FileKind,
}
/// The server's closing summary of a finished walk.
#[derive(Clone, Copy, PartialEq)]
struct Summary {
scanned: usize,
skipped: usize,
elapsed_ms: u64,
}
#[derive(Clone, Copy, PartialEq)]
enum Status {
Idle,
Searching,
/// The walk finished.
/// The walk ended. `stopped` marks an early end: the user pressed Stop,
/// or the server hit its match cap. `summary` is `None` when the user
/// stopped, because the connection is gone before the summary arrives.
Done {
scanned: usize,
skipped: usize,
elapsed_ms: u64,
stopped: bool,
summary: Option<Summary>,
},
/// The user pressed Stop. The connection was aborted before the server
/// could send its summary, so no scan counts are available.
StoppedByUser,
/// The server stopped early (match cap reached) and sent its summary.
Stopped {
scanned: usize,
skipped: usize,
},
}
/// A text fragment: plain or a highlighted hit. Rendered through
/// [`FragView`] so a mixed fragment list has one concrete element type
/// (Leptos 0.8 views are generic).
#[derive(Clone, PartialEq, Eq)]
enum Frag {
Plain(String),
Mark(String),
}
/// Render `text` with every query hit wrapped in `<mark class="hl-hit">`, as
@@ -255,21 +194,67 @@ enum Frag {
/// Both file names and file contents are untrusted, so every character of
/// `text` goes through [`push_escaped`] and the only markup in the output is
/// the `<mark>` this function writes.
///
/// Matching is case-insensitive: the text is lowercased into a copy with one
/// lowercase character per original character, so a hit's offsets in the
/// copy map straight back onto the original.
///
// ponytail: `lower` keeps that 1:1 map by taking only the first character of
// each `to_lowercase()`. A character whose lowercase expands (`İ`) therefore
// matches on its first character only. Use a real case-folding table if that
// ever matters.
fn highlight_html(text: &str, words: &[String]) -> String {
let mut out = String::with_capacity(text.len() + 16);
for f in highlight(text, words) {
match f {
Frag::Plain(s) => push_escaped(&mut out, &s),
Frag::Mark(s) => {
out.push_str("<mark class=\"hl-hit\">");
push_escaped(&mut out, &s);
out.push_str("</mark>");
}
let needles: Vec<String> = words
.iter()
.filter(|w| !w.is_empty())
.map(|w| lower(w))
.collect();
if needles.is_empty() || text.is_empty() {
push_escaped(&mut out, text);
return out;
}
// The lowered copy, plus the original byte offset behind every byte of
// it. `orig_at[i]` is meaningful at each character boundary `i`, and
// carries `text.len()` as its final sentinel.
let mut lowered = String::with_capacity(text.len());
let mut orig_at: Vec<usize> = Vec::with_capacity(text.len() + 1);
for (b, c) in text.char_indices() {
let lc = c.to_lowercase().next().unwrap_or(c);
lowered.push(lc);
for _ in 0..lc.len_utf8() {
orig_at.push(b);
}
}
orig_at.push(text.len());
let mut pos = 0usize;
while pos < lowered.len() {
// Earliest hit of any word at or after `pos`.
let Some((start, n)) = needles
.iter()
.filter_map(|w| lowered[pos..].find(w.as_str()).map(|i| (pos + i, w.len())))
.min_by_key(|(i, _)| *i)
else {
break;
};
push_escaped(&mut out, &text[orig_at[pos]..orig_at[start]]);
out.push_str("<mark class=\"hl-hit\">");
push_escaped(&mut out, &text[orig_at[start]..orig_at[start + n]]);
out.push_str("</mark>");
pos = start + n;
}
push_escaped(&mut out, &text[orig_at[pos]..]);
out
}
/// Lowercase, one character out per character in — see [`highlight_html`].
fn lower(s: &str) -> String {
s.chars()
.map(|c| c.to_lowercase().next().unwrap_or(c))
.collect()
}
/// Append `s` to `out` with the five HTML-significant characters escaped.
fn push_escaped(out: &mut String, s: &str) {
for c in s.chars() {
@@ -284,137 +269,24 @@ fn push_escaped(out: &mut String, s: &str) {
}
}
/// Split `text` into plain/highlighted fragments for the query's words
/// (case-insensitive). Name matches are whole-word ANDs and content matches
/// are the whole phrase, so highlighting each word covers both.
///
/// The text is lowercased once; hits are located in lowercased-char space
/// and mapped back to the original by character index, so case mappings
/// that change the character count (rare, e.g. `İ`) stay correct.
///
/// Worst case is O(words × len²) (the inner search restarts at `pos` for
/// every word); acceptable because it runs once per rendered row, not per
/// flush.
fn highlight(text: &str, words: &[String]) -> Vec<Frag> {
let needles: Vec<Vec<char>> = words
.iter()
.filter(|w| !w.is_empty())
.map(|w| w.to_lowercase().chars().collect())
.collect();
if needles.is_empty() {
return if text.is_empty() {
Vec::new()
} else {
vec![Frag::Plain(text.to_string())]
};
}
let lowered: Vec<char> = text.chars().flat_map(|c| c.to_lowercase()).collect();
// Byte offset just past original character `i`.
let mut byte_end: Vec<usize> = Vec::with_capacity(text.len());
for (b, ch) in text.char_indices() {
byte_end.push(b + ch.len_utf8());
}
// Lowered index -> original char index. Only differs from the identity
// when some character lowercases to more than one character.
let mut l2c: Vec<usize> = Vec::new();
if lowered.len() != byte_end.len() {
for ci in 0..byte_end.len() {
let s = if ci == 0 { 0 } else { byte_end[ci - 1] };
let ch = &text[s..byte_end[ci]];
for _ in ch.to_lowercase().chars() {
l2c.push(ci);
}
}
}
let char_of = |lp: usize| -> usize {
if l2c.is_empty() {
lp
} else {
l2c.get(lp).copied().unwrap_or(byte_end.len())
}
};
// Byte offset of original character `ci` (0 for the first).
let byte_at = |ci: usize| -> usize { if ci == 0 { 0 } else { byte_end[ci - 1] } };
let mut out: Vec<Frag> = Vec::new();
let mut pos = 0usize;
while pos < lowered.len() {
// Earliest hit of any word at or after `pos`.
let mut best: Option<(usize, usize)> = None;
for w in &needles {
let n = w.len();
if pos + n > lowered.len() {
continue;
}
let mut i = pos;
while i + n <= lowered.len() {
if lowered[i..i + n] == w[..] {
if best.is_none_or(|(b, _)| i < b) {
best = Some((i, n));
}
break;
}
i += 1;
}
}
let Some((start, n)) = best else {
break;
};
let s0 = byte_at(char_of(pos));
let s1 = byte_at(char_of(start));
let s2 = byte_at(char_of(start + n));
if s1 > s0 {
out.push(Frag::Plain(text[s0..s1].to_string()));
}
out.push(Frag::Mark(text[s1..s2].to_string()));
pos = start + n;
}
let st = byte_at(char_of(pos));
if st < text.len() {
out.push(Frag::Plain(text[st..].to_string()));
}
out
}
/// Parse `#/search?q=...&scope=...&root=...` (the part after `?`).
///
/// `roots` is still accepted, and its first id taken, so links made when a
/// search could span several roots still open.
fn parse_search_url() -> (String, Scope, Option<i64>) {
let hash = web_sys::window()
.and_then(|w| w.location().hash().ok())
.unwrap_or_default();
let Some(qs) = hash.split_once('?').map(|(_, q)| q) else {
let params = hash
.split_once('?')
.and_then(|(_, qs)| web_sys::UrlSearchParams::new_with_str(qs).ok());
let Some(p) = params else {
return (String::new(), Scope::Both, None);
};
let mut q = String::new();
let mut scope = Scope::Both;
let mut root = None;
for pair in qs.split('&') {
let (key, val) = match pair.split_once('=') {
Some(p) => p,
None => continue,
};
match key {
"q" => {
q = js_sys::decode_uri_component(val)
.map(String::from)
.unwrap_or_else(|_| val.to_string())
}
"scope" => {
scope = match val {
"name" => Scope::Name,
"content" => Scope::Content,
_ => Scope::Both,
}
}
"root" | "roots" => {
root = val.split(',').next().and_then(|s| s.trim().parse().ok());
}
_ => {}
}
}
(q, scope, root)
let scope = match p.get("scope").as_deref() {
Some("name") => Scope::Name,
Some("content") => Scope::Content,
_ => Scope::Both,
};
let root = p.get("root").and_then(|s| s.trim().parse().ok());
(p.get("q").unwrap_or_default(), scope, root)
}
/// Write the current search into the hash without adding a history entry.
@@ -455,34 +327,6 @@ fn goto_parent(root_id: i64, path: &str) {
});
}
/// The value of a `<select>` behind an event, when the target is one.
fn select_value(ev: &web_sys::Event) -> Option<String> {
ev.target()
.and_then(|t| t.dyn_into::<web_sys::HtmlSelectElement>().ok())
.map(|s| s.value())
}
/// Guess a file's kind from its name (search results carry no sniffed kind;
/// the browser listing does). Unknown extensions open in the "no preview"
/// view rather than the editor, so a misguess is a harmless dead end.
fn kind_of_name(name: &str) -> FileKind {
let ext = name.rsplit('.').next().unwrap_or("").to_lowercase();
match ext.as_str() {
"png" | "jpg" | "jpeg" | "gif" | "webp" | "bmp" | "svg" | "ico" => FileKind::Image,
"mp4" | "webm" | "mov" | "mkv" | "avi" => FileKind::Video,
"mp3" | "wav" | "ogg" | "flac" | "m4a" | "opus" => FileKind::Audio,
"pdf" => FileKind::Pdf,
"zip" | "tar" | "gz" | "tgz" | "bz2" | "xz" | "zst" | "7z" | "rar" => FileKind::Archive,
"txt" | "md" | "markdown" | "rs" | "py" | "js" | "mjs" | "ts" | "tsx" | "jsx" | "json"
| "jsonc" | "toml" | "yaml" | "yml" | "ini" | "cfg" | "conf" | "sh" | "bash" | "zsh"
| "c" | "h" | "cpp" | "hpp" | "cc" | "cs" | "go" | "java" | "kt" | "swift" | "rb"
| "php" | "html" | "htm" | "css" | "scss" | "less" | "xml" | "sql" | "csv" | "tsv"
| "log" | "tex" | "lua" | "pl" | "r" | "dart" | "vue" | "svelte" | "proto" | "graphql"
| "dockerfile" | "makefile" | "env" | "gitignore" => FileKind::Text,
_ => FileKind::Binary,
}
}
#[component]
pub fn SearchView(
me: ReadSignal<Option<api_types::Me>>,
@@ -495,7 +339,9 @@ pub fn SearchView(
// ---- state -------------------------------------------------------------
let (query, set_query) = RwSignal::<String>::new(String::new()).split();
let (scope, set_scope) = RwSignal::<Scope>::new(Scope::Both).split();
let (exec_query, set_exec_query) = RwSignal::<String>::new(String::new()).split();
// The running search's query words, for highlighting. Not reactive: rows
// read it while they render, and it only changes with a new search.
let words = StoredValue::new(Vec::<String>::new());
// The one root a search covers. Defaults to the caller's first root; the
// `<select>` re-reads `me` live for the names.
let (sel_root, set_sel_root) = RwSignal::<Option<i64>>::new(
@@ -527,28 +373,26 @@ pub fn SearchView(
// and reset on every new search.
let registry: RwSignal<HashMap<Arc<str>, CardState>> = RwSignal::new(HashMap::new());
// The in-flight AbortController, so Stop and unmount can both reach it.
// Deliberately *not* reactive: the fetch task outlives the component and
// The open event source, so Stop and unmount can both reach it.
// Deliberately *not* reactive: the stream outlives the component and
// touching a disposed signal panics.
let abort_box: Arc<Mutex<Option<web_sys::AbortController>>> = Arc::new(Mutex::new(None));
// Set on unmount so a stream event arriving one tick late cannot touch
// disposed signals.
let guard: Arc<AtomicBool> = Arc::new(AtomicBool::new(false));
// Stream events are buffered and applied in batches (see FLUSH_MS);
// `gen` is bumped on every new search so a stale flush from a previous
// run cannot leak its batch into the fresh state.
let source: Arc<Mutex<Option<web_sys::EventSource>>> = Arc::new(Mutex::new(None));
// Stream events are buffered and applied in batches (see FLUSH_MS).
// `search_gen` is bumped on every new search *and on unmount*, so one
// check covers both "a stale flush from a previous run" and "the view is
// gone, do not touch its signals".
let pending: Arc<Mutex<Vec<SearchEvent>>> = Arc::new(Mutex::new(Vec::new()));
let flushing: Arc<AtomicBool> = Arc::new(AtomicBool::new(false));
let search_gen: Arc<AtomicU64> = Arc::new(AtomicU64::new(0));
// ---- actions -----------------------------------------------------------
let stop_search = Callback::new({
let abort_box = abort_box.clone();
let source = source.clone();
let pending = pending.clone();
let search_gen = search_gen.clone();
move |_| {
if let Some(c) = abort_box.lock().unwrap().take() {
c.abort();
if let Some(s) = source.lock().unwrap().take() {
s.close();
}
// Stop should stop the UI too: discard buffered events and
// invalidate the in-flight flush so nothing keeps landing. The
@@ -556,9 +400,12 @@ pub fn SearchView(
// the "N more" rows report the difference to what is rendered.
pending.lock().unwrap().clear();
search_gen.fetch_add(1, Ordering::SeqCst);
// The server can no longer send its `Done` summary (the
// connection is gone), so settle the UI here.
status.set(Status::StoppedByUser);
// The server can no longer send its summary (the connection is
// gone), so settle the UI here.
status.set(Status::Done {
stopped: true,
summary: None,
});
}
});
@@ -568,7 +415,6 @@ pub fn SearchView(
let flush_ctx = FlushCtx {
pending: pending.clone(),
flushing: flushing.clone(),
guard: guard.clone(),
search_gen: search_gen.clone(),
files: files_sig,
files_total,
@@ -597,8 +443,8 @@ pub fn SearchView(
let start_search = Callback::new({
let stop = stop_search;
let abort_box = abort_box.clone();
let guard = guard.clone();
let source = source.clone();
let search_gen = search_gen.clone();
move |()| {
// Untracked reads: starting a search must never register reactive
// dependencies (this callback can run from any context).
@@ -608,9 +454,9 @@ pub fn SearchView(
return;
}
stop.run(()); // stop any in-flight search first
*abort_box.lock().unwrap() = None;
*source.lock().unwrap() = None;
pending.lock().unwrap().clear();
search_gen.fetch_add(1, Ordering::SeqCst);
let my_gen = search_gen.fetch_add(1, Ordering::SeqCst) + 1;
// Exactly one root. Falls back to the first available when
// nothing is selected yet (a deep link with an unknown id, or
@@ -621,7 +467,6 @@ pub fn SearchView(
}) else {
return;
};
let roots = vec![root];
set_sel_root.set(Some(root));
// Fresh state for this run.
@@ -636,36 +481,36 @@ pub fn SearchView(
registry.set(HashMap::new());
status.set(Status::Searching);
set_search_url(&q, scope, root);
set_exec_query.set(q.clone());
words.set_value(query_words(&q));
let ctrl = match api::search_stream(
let src = match api::search_stream(
q.clone(),
scope.param(),
&roots,
root,
on_stream_event,
on_error_for(toast, status, guard.clone()),
on_error_for(toast, status, search_gen.clone(), my_gen),
) {
Ok(c) => c,
Ok(s) => s,
Err(e) => {
show_error(toast, e.to_string());
status.set(Status::Idle);
return;
}
};
*abort_box.lock().unwrap() = Some(ctrl);
*source.lock().unwrap() = Some(src);
}
});
// The search view's owner is the shell's view tree; a stream event can
// arrive on the main thread after unmount, so the stream-side callbacks
// bail on the (non-reactive) guard, and the fetch is aborted here.
// arrive on the main thread after unmount, so bumping the generation here
// makes every stream-side callback bail, and the stream is closed.
{
let abort_box = abort_box.clone();
let g = guard.clone();
let source = source.clone();
let search_gen = search_gen.clone();
on_cleanup(move || {
g.store(true, Ordering::Relaxed);
if let Some(c) = abort_box.lock().unwrap().take() {
c.abort();
search_gen.fetch_add(1, Ordering::SeqCst);
if let Some(s) = source.lock().unwrap().take() {
s.close();
}
});
}
@@ -678,7 +523,8 @@ pub fn SearchView(
// signal while the registry's own borrow is held would rely on
// leptos deferring effects to a microtask.
let states: Vec<CardState> = registry.with_untracked(|m| m.values().copied().collect());
let target = !(!states.is_empty() && states.iter().all(|c| c.collapsed.get_untracked()));
let target =
!(!states.is_empty() && states.iter().all(|c| c.collapsed.get_untracked()));
for c in states {
c.collapsed.set(target);
}
@@ -742,7 +588,9 @@ pub fn SearchView(
// One delegated handler for the whole name list, instead of a `Callback`
// per row: a callback is an arena entry plus a closure allocation, and at
// a thousand rows that is a measurable slice of the render. The row
// carries what the handler needs in data attributes.
// carries its identity in data attributes; the hit itself is looked up in
// the list, which is a linear scan of at most `RETAIN_FILE_CAP` on a
// click.
let on_row_click = move |ev: web_sys::MouseEvent| {
let Some(target) = ev
.target()
@@ -768,7 +616,14 @@ pub fn SearchView(
goto_parent(root_id, &path);
return;
}
if row.get_attribute("data-dir").is_some() {
let Some(hit) = files_r.with_untracked(|v| {
v.iter()
.find(|h| h.root_id == root_id && *h.path == *path)
.cloned()
}) else {
return;
};
if hit.is_dir {
navigate(&Location {
root_id: Some(root_id),
path: split_rel(&path),
@@ -777,20 +632,15 @@ pub fn SearchView(
});
return;
}
let is_rw = me
.get_untracked()
.map(|m| {
m.roots
.iter()
.find(|r| r.id == root_id)
.map(|r| r.mode.is_writable())
.unwrap_or(false)
})
.unwrap_or(false);
let name = path.rsplit('/').next().unwrap_or(&path).to_string();
let kind = kind_of_name(&name);
saved_scroll.set_value(scroll_y());
open_file.set(Some(open_file_view(root_id, path, name, kind, is_rw)));
open_file.set(Some(open_file_view(
root_id,
path,
name,
hit.kind,
is_writable(me, root_id),
)));
};
// ---- render ------------------------------------------------------------
@@ -848,7 +698,11 @@ pub fn SearchView(
<select
class="root-select"
on:change=move |ev: web_sys::Event| {
if let Some(id) = select_value(&ev).and_then(|v| v.parse::<i64>().ok()) {
if let Some(id) = ev
.target()
.and_then(|t| t.dyn_into::<web_sys::HtmlSelectElement>().ok())
.and_then(|s| s.value().parse::<i64>().ok())
{
set_sel_root.set(Some(id));
}
}
@@ -891,7 +745,7 @@ pub fn SearchView(
}}
</div>
{move || match status.get() {
{move || match status.get() {
Status::Idle => view! {}.into_view().into_any(),
Status::Searching => view! {
<div class="status-line">
@@ -908,70 +762,46 @@ pub fn SearchView(
}
.into_view()
.into_any(),
Status::Done { scanned, skipped, elapsed_ms } => view! {
<div class="status-line finished">
<span>
{i18n::t_fmt(k::SEARCH_N_RESULTS, &(files_total.get() + matches_total.get()).to_string())}
{" · "}
{i18n::t_fmt(k::SEARCH_TOOK, &elapsed_ms.to_string())}
{" · "}
{i18n::t_fmt(k::SEARCH_SCANNED, &scanned.to_string())}
{move || {
if skipped > 0 {
view! {
<span class="skip">
{" · "}
{i18n::t_fmt(k::SEARCH_SKIPPED, &skipped.to_string())}
</span>
}
.into_view()
.into_any()
} else {
view! {}.into_view().into_any()
}
}}
</span>
</div>
}
.into_view()
.into_any(),
Status::StoppedByUser => view! {
<div class="status-line finished">
<span>
{i18n::t_fmt(
k::SEARCH_STOPPED,
&(files_total.get() + matches_total.get()).to_string(),
)}
</span>
</div>
}
.into_view()
.into_any(),
Status::Stopped { scanned, skipped } => view! {
<div class="status-line finished">
<span>
{i18n::t_fmt(k::SEARCH_STOPPED, &(files_total.get() + matches_total.get()).to_string())}
{" · "}
{i18n::t_fmt(k::SEARCH_SCANNED, &scanned.to_string())}
{move || {
if skipped > 0 {
view! {
<span class="skip">
{" · "}
{i18n::t_fmt(k::SEARCH_SKIPPED, &skipped.to_string())}
</span>
}
.into_view()
.into_any()
} else {
view! {}.into_view().into_any()
}
}}
</span>
</div>
Status::Done { stopped, summary } => {
let n = (files_total.get() + matches_total.get()).to_string();
let head = if stopped {
i18n::t_fmt(k::SEARCH_STOPPED, &n)
} else {
i18n::t_fmt(k::SEARCH_N_RESULTS, &n)
};
view! {
<div class="status-line finished">
<span>
{head}
{summary
.map(|s| {
view! {
<>
{" · "}
{i18n::t_fmt(k::SEARCH_TOOK, &s.elapsed_ms.to_string())}
{" · "}
{i18n::t_fmt(k::SEARCH_SCANNED, &s.scanned.to_string())}
{(s.skipped > 0)
.then(|| {
view! {
<span class="skip">
{" · "}
{i18n::t_fmt(
k::SEARCH_SKIPPED,
&s.skipped.to_string(),
)}
</span>
}
})}
</>
}
})}
</span>
</div>
}
.into_view()
.into_any()
}
.into_view()
.into_any(),
}}
// The groups are hidden with `display` instead of unmounted:
@@ -1008,17 +838,16 @@ pub fn SearchView(
// Note: only the first `shown_files` hits reach
// here; `each` above slices before cloning.
let (dir, name) = split_name(&h.path);
let words = query_words(&exec_query.get_untracked());
let icon = if h.is_dir {
IconName::Folder
} else {
icon_for(kind_of_name(name), name)
icon_for(h.kind, name)
};
// Only the name is highlighted. The server matches
// the entry's own name, not its path, so a mark
// in the directory line would point at text that
// was never matched.
let name_h = highlight_html(name, &words);
let name_h = words.with_value(|w| highlight_html(name, w));
let dir_s = dir.to_string();
let size = if h.is_dir {
"—".to_string()
@@ -1034,7 +863,6 @@ pub fn SearchView(
class="srow"
data-root=h.root_id.to_string()
data-path=h.path.to_string()
data-dir=h.is_dir.then_some("")
>
{icon_svg(icon, "srow-icon")}
<div class="srow-main">
@@ -1139,7 +967,10 @@ pub fn SearchView(
let key = f.path.clone();
let root_id = f.root_id;
let (dir, name) = split_name(&f.path);
let icon = icon_for(kind_of_name(name), name);
// A content match came out of the grep, which only
// reads text files, so the kind is known without
// asking the server.
let icon = icon_for(FileKind::Text, name);
// The card's path is *not* highlighted: a content hit
// was found in the file's text, not in its name, so
// marking the name would claim a match that is not
@@ -1213,9 +1044,8 @@ pub fn SearchView(
key=move |p: &(usize, MatchLine)| (p.1.line, p.0)
children=move |p: (usize, MatchLine)| {
let l = p.1;
let words =
query_words(&exec_query.get_untracked());
let text_h = highlight_html(&l.text, &words);
let text_h = words
.with_value(|w| highlight_html(&l.text, w));
view! {
<div class="mline">
<span class="ln">{l.line}</span>
@@ -1244,10 +1074,7 @@ pub fn SearchView(
</div>
{move || {
let done = matches!(
status.get(),
Status::Done { .. } | Status::Stopped { .. } | Status::StoppedByUser
);
let done = matches!(status.get(), Status::Done { .. });
if done && files_total.get() == 0 && matches_total.get() == 0 {
view! {
<div class="search-empty">{i18n::t(k::SEARCH_NO_RESULTS)}</div>
@@ -1267,7 +1094,6 @@ pub fn SearchView(
struct FlushCtx {
pending: Arc<Mutex<Vec<SearchEvent>>>,
flushing: Arc<AtomicBool>,
guard: Arc<AtomicBool>,
search_gen: Arc<AtomicU64>,
files: RwSignal<Vec<FileHit>>,
files_total: RwSignal<usize>,
@@ -1279,36 +1105,36 @@ struct FlushCtx {
status: RwSignal<Status>,
}
impl FlushCtx {
/// True once a new search has started or the view was unmounted — both
/// bump the generation, and after either the view signals must not be
/// written any more.
fn stale(&self, my_gen: u64) -> bool {
self.search_gen.load(Ordering::SeqCst) != my_gen
}
}
/// Waits `delay_ms`, then applies every event buffered in the meantime.
///
/// `flushing` is held for the whole (yielding) apply, so only one flush ever
/// writes the view state; events arriving during it buffer in `pending` and
/// the next flush picks them up right away instead of waiting `FLUSH_MS`.
async fn run_flush(delay_ms: u32, my_gen: u64, ctx: FlushCtx) {
let _ = TimeoutFuture::new(delay_ms).await;
let stale =
|| ctx.guard.load(Ordering::Relaxed) || ctx.search_gen.load(Ordering::SeqCst) != my_gen;
if stale() {
ctx.flushing.store(false, Ordering::SeqCst);
return;
}
let batch: Vec<SearchEvent> = ctx.pending.lock().unwrap().drain(..).collect();
if !batch.is_empty() {
apply_batch(&ctx, batch, my_gen).await;
TimeoutFuture::new(delay_ms).await;
if !ctx.stale(my_gen) {
let batch: Vec<SearchEvent> = ctx.pending.lock().unwrap().drain(..).collect();
if !batch.is_empty() {
apply_batch(&ctx, batch, my_gen).await;
}
// Applying yields, so more events may have arrived meanwhile. They
// are already `FLUSH_MS` old, so pick them up on the next tick —
// keeping `flushing` held, so no second flush can start.
if !ctx.stale(my_gen) && !ctx.pending.lock().unwrap().is_empty() {
wasm_bindgen_futures::spawn_local(run_flush(0, my_gen, ctx.clone()));
return;
}
}
ctx.flushing.store(false, Ordering::SeqCst);
// Applying yields, so more events may have arrived meanwhile. They are
// already `FLUSH_MS` old, so pick them up on the next frame.
if stale() || ctx.pending.lock().unwrap().is_empty() {
return;
}
if ctx
.flushing
.compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
.is_ok()
{
wasm_bindgen_futures::spawn_local(run_flush(0, my_gen, ctx.clone()));
}
}
/// Folds one batch into the view state with a single render per signal:
@@ -1337,7 +1163,7 @@ async fn apply_batch(ctx: &FlushCtx, batch: Vec<SearchEvent>, my_gen: u64) {
let mut new_match_files: Vec<MatchFile> = Vec::new();
let mut mf_len = ctx.match_files.with_untracked(|v| v.len());
// Generation stamp for the card keys (see `MatchFile::search_gen`).
let search_gen = ctx.search_gen.load(Ordering::SeqCst);
let search_gen = my_gen;
let mut done: Option<Status> = None;
// Group the batch's match events per file, so a new card is seeded with
@@ -1358,9 +1184,7 @@ async fn apply_batch(ctx: &FlushCtx, batch: Vec<SearchEvent>, my_gen: u64) {
if since_yield >= EVENTS_PER_SLICE {
since_yield = 0;
yield_to_browser().await;
if ctx.guard.load(Ordering::Relaxed)
|| ctx.search_gen.load(Ordering::SeqCst) != my_gen
{
if ctx.stale(my_gen) {
return;
}
}
@@ -1373,14 +1197,13 @@ async fn apply_batch(ctx: &FlushCtx, batch: Vec<SearchEvent>, my_gen: u64) {
skipped,
elapsed_ms,
} => {
done = Some(if stopped {
Status::Stopped { scanned, skipped }
} else {
Status::Done {
done = Some(Status::Done {
stopped,
summary: Some(Summary {
scanned,
skipped,
elapsed_ms,
}
}),
});
}
SearchEvent::File {
@@ -1388,6 +1211,7 @@ async fn apply_batch(ctx: &FlushCtx, batch: Vec<SearchEvent>, my_gen: u64) {
path,
size,
is_dir,
kind,
} => {
file_count += 1;
if file_budget > 0 {
@@ -1397,6 +1221,7 @@ async fn apply_batch(ctx: &FlushCtx, batch: Vec<SearchEvent>, my_gen: u64) {
path: path.into(),
size,
is_dir,
kind,
});
}
}
@@ -1509,13 +1334,16 @@ async fn apply_batch(ctx: &FlushCtx, batch: Vec<SearchEvent>, my_gen: u64) {
.await;
}
/// Reports a stream failure, unless the search it belongs to is already gone
/// (a new search, or the view unmounted — both bump the generation).
fn on_error_for(
toast: ToastMsg,
status: RwSignal<Status>,
guard: Arc<AtomicBool>,
search_gen: Arc<AtomicU64>,
my_gen: u64,
) -> Callback<String, ()> {
Callback::new(move |msg| {
if guard.load(Ordering::Relaxed) {
if search_gen.load(Ordering::SeqCst) != my_gen {
return;
}
show_error(toast, msg);
@@ -1523,6 +1351,16 @@ fn on_error_for(
})
}
/// Whether the signed-in user may write to `root_id`.
fn is_writable(me: ReadSignal<Option<api_types::Me>>, root_id: i64) -> bool {
me.get_untracked().is_some_and(|m| {
m.roots
.iter()
.find(|r| r.id == root_id)
.is_some_and(|r| r.mode.is_writable())
})
}
fn open_file_view(
root_id: i64,
path: String,
@@ -1556,11 +1394,11 @@ fn open_file_view(
}
}
/// Open a content match: always a file. One callback per card, shared by all
/// of its lines. Free function rather than a component closure so the card's
/// `<For>` children closure only captures `Copy` signal handles (the
/// surrounding closures must be re-runnable, and `move` closures may only
/// move `Copy` captures out of that environment).
/// Open a content match: always a text file (the grep only reads those). One
/// callback per card, shared by all of its lines. Free function rather than a
/// component closure so the card's `<For>` children closure only captures
/// `Copy` signal handles (the surrounding closures must be re-runnable, and
/// `move` closures may only move `Copy` captures out of that environment).
fn open_match_cb(
root_id: i64,
path: Arc<str>,
@@ -1568,19 +1406,14 @@ fn open_match_cb(
open_file: WriteSignal<Option<FileView>>,
) -> Callback<()> {
Callback::new(move |_| {
let is_rw = me
.get_untracked()
.map(|m| {
m.roots
.iter()
.find(|r| r.id == root_id)
.map(|r| r.mode.is_writable())
.unwrap_or(false)
})
.unwrap_or(false);
let name = path.rsplit('/').next().unwrap_or(&path).to_string();
let kind = kind_of_name(&name);
let fv = open_file_view(root_id, path.to_string(), name, kind, is_rw);
let fv = open_file_view(
root_id,
path.to_string(),
name,
FileKind::Text,
is_writable(me, root_id),
);
open_file.set(Some(fv));
})
}
@@ -1627,7 +1460,7 @@ fn page_footer(
/// Lowercased whitespace-split words of the query.
fn query_words(query: &str) -> Vec<String> {
query.split_whitespace().map(|w| w.to_lowercase()).collect()
query.split_whitespace().map(lower).collect()
}
fn split_rel(path: &str) -> Vec<String> {
@@ -1649,89 +1482,65 @@ fn split_name(path: &str) -> (&str, &str) {
mod tests {
use super::*;
fn frags(f: &[Frag]) -> Vec<String> {
f.iter()
.map(|f| match f {
Frag::Plain(s) => format!("={s}="),
Frag::Mark(s) => format!("<{s}>"),
})
.collect()
fn hl(text: &str, words: &[&str]) -> String {
let w: Vec<String> = words.iter().map(|s| s.to_string()).collect();
highlight_html(text, &w)
}
/// `<mark class="hl-hit">x</mark>` written as `<x>`, so the expectations
/// stay readable.
fn short(html: &str) -> String {
html.replace("<mark class=\"hl-hit\">", "<")
.replace("</mark>", ">")
}
#[test]
fn highlight_marks_hits() {
let w = vec!["world".to_string()];
assert_eq!(
frags(&highlight("Hello world", &w)),
vec!["=Hello =", "<world>"]
);
assert_eq!(short(&hl("Hello world", &["world"])), "Hello <world>");
}
#[test]
fn highlight_is_case_insensitive() {
let w = vec!["hello".to_string()];
assert_eq!(
frags(&highlight("HeLLo there", &w)),
vec!["<HeLLo>", "= there="]
);
assert_eq!(short(&hl("HeLLo there", &["hello"])), "<HeLLo> there");
}
#[test]
fn highlight_marks_all_words() {
let w = vec!["foo".to_string(), "bar".to_string()];
assert_eq!(
frags(&highlight("foo bar foo", &w)),
vec!["<foo>", "= =", "<bar>", "= =", "<foo>"]
short(&hl("foo bar foo", &["foo", "bar"])),
"<foo> <bar> <foo>"
);
}
#[test]
fn highlight_no_hit_is_plain() {
let w = vec!["xyz".to_string()];
assert_eq!(frags(&highlight("abc", &w)), vec!["=abc="]);
assert_eq!(hl("abc", &["xyz"]), "abc");
}
#[test]
fn highlight_empty_words_is_plain() {
assert_eq!(frags(&highlight("abc", &[])), vec!["=abc="]);
assert!(highlight("", &["a".to_string()]).is_empty());
assert_eq!(hl("abc", &[]), "abc");
assert_eq!(hl("", &["a"]), "");
}
#[test]
fn highlight_html_escapes_markup() {
fn highlight_escapes_markup() {
// File names and file contents are untrusted: the only markup in the
// output must be the `<mark>` wrapper.
let w = vec!["script".to_string()];
assert_eq!(
highlight_html("<script>alert(1)</script>", &w),
hl("<script>alert(1)</script>", &["script"]),
"<<mark class=\"hl-hit\">script</mark>>alert(1)</\
<mark class=\"hl-hit\">script</mark>>"
);
// Escaping applies inside the highlighted span too.
assert_eq!(
highlight_html("a<b", &["<b".to_string()]),
"a<mark class=\"hl-hit\"><b</mark>"
);
assert_eq!(
highlight_html(r#"& " ' < >"#, &[]),
"& " ' < >"
);
assert_eq!(hl("a<b", &["<b"]), "a<mark class=\"hl-hit\"><b</mark>");
assert_eq!(hl(r#"& " ' < >"#, &[]), "& " ' < >");
}
#[test]
fn highlight_handles_non_ascii() {
// Ö (2 bytes) and ö (2 bytes) vs. the ASCII around them: the hit
// slices must follow character boundaries of the original text.
let w = vec!["ö".to_string()];
let out = highlight("AÖ aö b", &w);
assert_eq!(frags(&out), vec!["=A=", "<Ö>", "= a=", "<ö>", "= b="]);
// The plain and mark fragments must reassemble the input.
let joined: String = out
.iter()
.map(|f| match f {
Frag::Plain(s) | Frag::Mark(s) => s.as_str(),
})
.collect();
assert_eq!(joined, "AÖ aö b");
// Ö and ö are 2 bytes each: the hit slices must follow character
// boundaries of the original text.
assert_eq!(short(&hl("AÖ aö b", &["ö"])), "A<Ö> a<ö> b");
}
}