#![allow(
clippy::cast_possible_truncation,
clippy::cast_possible_wrap,
clippy::cast_precision_loss,
clippy::cast_sign_loss
)]
use std::fmt::Write as _;
use std::path::Path;
use crate::db::FileRow;
use anyhow::{Context, Result};
use sha2::{Digest, Sha256};
use super::App;
pub enum OcrUpdate {
Stage(String),
Progress(usize, usize, usize),
Done(std::result::Result<(String, Vec<(usize, String)>), String>),
}
pub struct PickerEntry {
pub name: String,
pub is_dir: bool,
}
#[derive(Clone)]
pub enum OcrBackend {
Router(crate::provider::openrouter::OpenRouter, String),
Ollama(reqwest::Client, String),
}
impl OcrBackend {
async fn transcribe(&self, png: &[u8]) -> anyhow::Result<String> {
self.transcribe_image(png, "image/png").await
}
async fn describe(&self, bytes: &[u8], mime: &str) -> anyhow::Result<String> {
match self {
Self::Router(provider, model) => {
use base64::Engine;
let b64 = base64::engine::general_purpose::STANDARD.encode(bytes);
let url = format!("data:{mime};base64,{b64}");
provider.describe_image(model, &url).await
}
Self::Ollama(client, model) => {
use base64::Engine;
let b64 = base64::engine::general_purpose::STANDARD.encode(bytes);
let resp = client
.post("http://127.0.0.1:11434/api/generate")
.timeout(std::time::Duration::from_mins(10))
.json(&serde_json::json!({
"model": model,
"prompt": "Describe this image so another AI model can reason about \
it without seeing it. Cover: what it is (screenshot, chart, \
photo, diagram…), overall layout and structure, the key \
entities and how they relate, ALL visible text verbatim, \
and any notable visual details. Be thorough but do not \
speculate beyond what is visible.",
"images": [b64],
"stream": false,
"options": { "num_ctx": 8192 },
}))
.send()
.await
.map_err(|e| {
if e.is_timeout() {
anyhow::anyhow!("timeout after 600s")
} else if e.is_connect() {
anyhow::anyhow!("cannot reach ollama — is it running?")
} else {
e.into()
}
})?;
if resp.status().as_u16() == 404 {
anyhow::bail!("model '{model}' not pulled");
}
let v = resp.error_for_status()?.json::<serde_json::Value>().await?;
Ok(v.get("response")
.and_then(|r| r.as_str())
.unwrap_or("")
.to_string())
}
}
}
async fn transcribe_image(&self, bytes: &[u8], mime: &str) -> anyhow::Result<String> {
match self {
Self::Router(provider, model) => {
use base64::Engine;
let b64 = base64::engine::general_purpose::STANDARD.encode(bytes);
let url = format!("data:{mime};base64,{b64}");
provider.ocr_page(model, &url).await
}
Self::Ollama(client, model) => {
use base64::Engine;
let b64 = base64::engine::general_purpose::STANDARD.encode(bytes);
let resp = client
.post("http://127.0.0.1:11434/api/generate")
.timeout(std::time::Duration::from_mins(10))
.json(&ollama_ocr_body(model, &b64))
.send()
.await
.map_err(|e| {
if e.is_timeout() {
anyhow::anyhow!("timeout after 600s")
} else if e.is_connect() {
anyhow::anyhow!(
"cannot reach ollama at 127.0.0.1:11434 — is it running? (systemctl start ollama)"
)
} else {
e.into()
}
})?;
if resp.status().as_u16() == 404 {
anyhow::bail!(
"model '{model}' not pulled — cycle OCR engine to 'local' in /config"
);
}
let v = resp.error_for_status()?.json::<serde_json::Value>().await?;
Ok(v.get("response")
.and_then(|r| r.as_str())
.unwrap_or("")
.to_string())
}
}
}
}
fn is_uuid_like(stem: &str) -> bool {
let hex = |s: &str| s.chars().all(|c| c.is_ascii_hexdigit());
let parts: Vec<&str> = stem.split('-').collect();
parts.len() == 5
&& parts[0].len() == 8
&& hex(parts[0])
&& parts[1].len() == 4
&& hex(parts[1])
&& parts[2].len() == 4
&& hex(parts[2])
&& parts[3].len() == 4
&& hex(parts[3])
&& parts[4].len() == 12
&& hex(parts[4])
}
fn clip_err(e: &str) -> String {
let mut s: String = e.chars().take(90).collect();
if s.len() < e.len() {
s.push('…');
}
s
}
fn ollama_ocr_body(model: &str, png_b64: &str) -> serde_json::Value {
serde_json::json!({
"model": model,
"prompt": crate::provider::openrouter::OCR_PROMPT,
"images": [png_b64],
"stream": false,
"options": { "num_ctx": 8192 },
})
}
async fn ocr_pdf_vlm(
backend: &OcrBackend,
path: &Path,
tx: &tokio::sync::mpsc::UnboundedSender<(String, String, OcrUpdate)>,
space_id: &str,
name: &str,
files_dir: &Path,
) -> std::result::Result<(String, Vec<(usize, String)>), String> {
let stem = std::path::Path::new(name)
.file_stem()
.and_then(|s| s.to_str())
.unwrap_or(name);
let page_dir = files_dir.join(stem);
if let Err(e) = std::fs::create_dir_all(&page_dir) {
return Err(format!("error: ocr: {e}"));
}
ocr_pdf_vlm_in(backend, path, &page_dir, tx, space_id, name).await
}
async fn ocr_pdf_vlm_in(
backend: &OcrBackend,
path: &Path,
page_dir: &Path,
tx: &tokio::sync::mpsc::UnboundedSender<(String, String, OcrUpdate)>,
space_id: &str,
name: &str,
) -> std::result::Result<(String, Vec<(usize, String)>), String> {
let _ = tx.send((
space_id.to_string(),
name.to_string(),
OcrUpdate::Stage("rendering pages (300 dpi)…".to_string()),
));
let (pdf, dir) = (path.to_path_buf(), page_dir.to_path_buf());
let pages = tokio::task::spawn_blocking(move || {
crate::extract::render_pdf_pages("pdftoppm", &pdf, &dir, 300, false)
})
.await
.map_err(|e| format!("error: ocr: {e}"))?
.map_err(|e| match e {
crate::extract::OcrError::MissingTools => {
"scanned pdf — install poppler (pdftoppm) for ocr".to_string()
}
crate::extract::OcrError::Failed(m) => format!("error: ocr: {m}"),
})?;
let total = pages.len();
let _ = tx.send((
space_id.to_string(),
name.to_string(),
OcrUpdate::Progress(0, total, 0),
));
let mut results: Vec<std::result::Result<String, String>> =
vec![Err("not transcribed".to_string()); total];
let mut set = tokio::task::JoinSet::new();
let spawn_page =
|set: &mut tokio::task::JoinSet<(usize, std::result::Result<String, String>)>, i: usize| {
let (backend, png) = (backend.clone(), pages[i].clone());
set.spawn(async move {
let Ok(bytes) = std::fs::read(&png) else {
return (i, Err("page image unreadable".to_string()));
};
let mut last = String::new();
for _ in 0..2 {
match backend.transcribe(&bytes).await {
Ok(text) => return (i, Ok(text)),
Err(e) => last = e.to_string(),
}
}
(i, Err(last))
});
};
let window = (16_usize).min(total);
let mut next = 0;
while next < window {
spawn_page(&mut set, next);
next += 1;
}
let mut done = 0;
let mut failed = 0;
while let Some(joined) = set.join_next().await {
let (i, r) = joined.unwrap_or_else(|_| (usize::MAX, Err("page task panicked".to_string())));
if r.is_err() {
failed += 1;
}
if let Some(slot) = results.get_mut(i) {
*slot = r;
}
done += 1;
let _ = tx.send((
space_id.to_string(),
name.to_string(),
OcrUpdate::Progress(done, total, failed),
));
if next < total {
spawn_page(&mut set, next);
next += 1;
}
}
let _ = std::fs::create_dir_all(page_dir);
for (i, p) in pages.iter().enumerate() {
let stable = page_dir.join(format!("page-{}.png", i + 1));
let _ = std::fs::rename(p, &stable);
}
let errors: Vec<(usize, String)> = results
.iter()
.enumerate()
.filter_map(|(i, r)| r.as_ref().err().map(|e| (i, e.clone())))
.collect();
Ok((crate::extract::join_pages(&results), errors))
}
async fn ocr_image_vlm(
backend: &OcrBackend,
path: &Path,
tx: &tokio::sync::mpsc::UnboundedSender<(String, String, OcrUpdate)>,
space_id: &str,
name: &str,
) -> std::result::Result<(String, Vec<(usize, String)>), String> {
let _ = tx.send((
space_id.to_string(),
name.to_string(),
OcrUpdate::Stage("transcribing image…".to_string()),
));
let Ok(bytes) = std::fs::read(path) else {
return Err(format!("cannot read {name}"));
};
let ext = path
.extension()
.and_then(|e| e.to_str())
.unwrap_or("")
.to_lowercase();
let mime = match ext.as_str() {
"jpg" | "jpeg" => "image/jpeg",
"gif" => "image/gif",
"webp" => "image/webp",
"bmp" => "image/bmp",
_ => "image/png",
};
let _ = tx.send((
space_id.to_string(),
name.to_string(),
OcrUpdate::Progress(0, 1, 0),
));
match backend.describe(&bytes, mime).await {
Ok(text) => {
let _ = tx.send((
space_id.to_string(),
name.to_string(),
OcrUpdate::Progress(1, 1, 0),
));
Ok((text, Vec::new()))
}
Err(e) => {
let err = e.to_string();
Err(format!("error: ocr: {err}"))
}
}
}
impl App {
pub(crate) fn open_file_picker(&mut self) {
self.picker_filter.clear();
self.picker_selected = 0;
self.reload_picker_entries();
self.files_mode = super::FilesMode::Pick;
}
fn reload_picker_entries(&mut self) {
let mut entries: Vec<PickerEntry> = match std::fs::read_dir(&self.picker_dir) {
Ok(rd) => rd
.flatten()
.filter_map(|e| {
let name = e.file_name().to_string_lossy().to_string();
let is_dir = e.file_type().ok()?.is_dir();
Some(PickerEntry { name, is_dir })
})
.collect(),
Err(e) => {
self.status = format!("cannot read {}: {e}", self.picker_dir.display());
Vec::new()
}
};
entries.sort_by(|a, b| b.is_dir.cmp(&a.is_dir).then_with(|| a.name.cmp(&b.name)));
self.picker_entries = entries;
}
pub fn filtered_picker_entries(&self) -> Vec<&PickerEntry> {
use crate::input::fuzzy_score;
let needle = self.picker_filter.trim();
if needle.is_empty() {
return self.picker_entries.iter().collect();
}
super::fuzzy_filter_sorted(&self.picker_entries, |e| fuzzy_score(&e.name, needle))
}
pub fn move_picker_selection(&mut self, delta: i32) {
self.picker_selected = super::clamp_cursor(
self.picker_selected,
self.filtered_picker_entries().len(),
delta,
);
}
pub fn picker_filter_push(&mut self, c: char) {
self.picker_filter.push(c);
self.picker_selected = 0;
}
pub fn picker_backspace(&mut self) {
if !self.picker_filter.is_empty() {
self.picker_filter.pop();
self.picker_selected = 0;
return;
}
if let Some(parent) = self.picker_dir.parent().map(std::path::Path::to_path_buf) {
self.picker_dir = parent;
self.picker_selected = 0;
self.reload_picker_entries();
}
}
pub fn picker_enter(&mut self) {
let filtered = self.filtered_picker_entries();
let Some(entry) = filtered.get(self.picker_selected) else {
return;
};
let name = entry.name.clone();
let is_dir = entry.is_dir;
let path = self.picker_dir.join(&name);
if is_dir {
self.picker_dir = path;
self.picker_filter.clear();
self.picker_selected = 0;
self.reload_picker_entries();
return;
}
match self.import_file(&path) {
Ok(n) => self.status = format!("imported {n}"),
Err(e) => self.status = format!("import failed: {e}"),
}
self.files_mode = super::FilesMode::Browse;
}
pub fn rescan_files(&mut self) {
let dir = self.space.files_dir(&self.active_space.name);
let known = self
.db
.list_files(&self.active_space.id)
.unwrap_or_default();
let mut seen: Vec<String> = Vec::new();
let mut ocr_jobs: Vec<(String, String, std::path::PathBuf)> = Vec::new();
let entries = std::fs::read_dir(&dir)
.map(|rd| rd.flatten().collect::<Vec<_>>())
.unwrap_or_default();
for entry in entries {
let path = entry.path();
if !path.is_file() {
continue;
}
let name = entry.file_name().to_string_lossy().to_string();
seen.push(name.clone());
let mtime = entry
.metadata()
.ok()
.and_then(|m| m.modified().ok())
.and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
.map_or(0, |d| d.as_secs() as i64);
let disk_size = entry.metadata().map_or(0, |m| m.len() as i64);
let existing = known.iter().find(|f| f.name == name);
if let Some(f) = existing
&& f.size == disk_size
&& f.mtime == mtime
&& mtime != 0
{
if f.status.starts_with("ocr") && self.ocr_rx.is_none() {
ocr_jobs.push((self.active_space.id.clone(), name.clone(), path.clone()));
}
continue;
}
let Ok(bytes) = std::fs::read(&path) else {
continue;
};
let hash = Sha256::digest(&bytes)
.iter()
.fold(String::new(), |mut h, b| {
let _ = write!(h, "{b:02x}");
h
});
if let Some(f) = existing.filter(|f| f.hash == hash) {
let _ = self.db.set_file_mtime(&f.id, mtime);
if f.status.starts_with("ocr") && self.ocr_rx.is_none() {
ocr_jobs.push((self.active_space.id.clone(), name.clone(), path.clone()));
}
continue;
}
let size = bytes.len() as i64;
let (status, chunks) = match crate::extract::extract_text(&path) {
Ok(text) if text.trim().is_empty() => {
let ext = std::path::Path::new(&name)
.extension()
.and_then(|e| e.to_str())
.unwrap_or("")
.to_lowercase();
if ext == "pdf" || crate::extract::is_image_ext(&ext) {
ocr_jobs.push((self.active_space.id.clone(), name.clone(), path.clone()));
("ocr…".to_string(), Vec::new())
} else {
("no text (scanned?)".to_string(), Vec::new())
}
}
Ok(text) => ("ok".to_string(), crate::extract::chunk_lines(&text)),
Err(e) => (format!("error: {e}"), Vec::new()),
};
if let Ok(id) = self
.db
.upsert_file(&self.active_space.id, &name, &hash, size, &status)
{
let _ = self.db.set_file_chunks(&id, &chunks);
let _ = self.db.set_file_mtime(&id, mtime);
}
}
for gone in known.iter().filter(|f| !seen.contains(&f.name)) {
let _ = self.db.delete_file(&gone.id);
}
self.start_ocr(ocr_jobs);
self.start_embedding();
self.files_cache = self
.db
.list_files(&self.active_space.id)
.unwrap_or_default();
self.files_selected = self
.files_selected
.min(self.files_cache.len().saturating_sub(1));
}
pub(crate) fn start_ocr(&mut self, jobs: Vec<(String, String, std::path::PathBuf)>) {
if jobs.is_empty() || self.ocr_rx.is_some() {
return;
}
let backend = self.ocr_backend();
let files_dir = self.space.files_dir(&self.active_space.name);
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
self.ocr_rx = Some(rx);
if let Some(backend) = backend {
tokio::spawn(async move {
for (space_id, name, path) in jobs {
let is_image = path
.extension()
.and_then(|e| e.to_str())
.is_some_and(crate::extract::is_image_ext);
let result = if is_image {
ocr_image_vlm(&backend, &path, &tx, &space_id, &name).await
} else {
ocr_pdf_vlm(&backend, &path, &tx, &space_id, &name, &files_dir).await
};
if tx.send((space_id, name, OcrUpdate::Done(result))).is_err() {
return;
}
}
});
return;
}
tokio::task::spawn_blocking(move || {
for (space_id, name, path) in jobs {
let is_image = path
.extension()
.and_then(|e| e.to_str())
.is_some_and(crate::extract::is_image_ext);
if is_image {
let _ = tx.send((
space_id,
name,
OcrUpdate::Done(Err("no vlm backend for image ocr".to_string())),
));
continue;
}
let progress_tx = tx.clone();
let (sid, fname) = (space_id.clone(), name.clone());
let progress = move |done: usize, total: usize| {
let _ = progress_tx.send((
sid.clone(),
fname.clone(),
OcrUpdate::Progress(done, total, 0),
));
};
let result = match crate::extract::ocr_pdf(&path, &progress) {
Ok(text) => Ok((text, Vec::new())),
Err(crate::extract::OcrError::MissingTools) => {
Err("scanned pdf — install tesseract + poppler for ocr".to_string())
}
Err(crate::extract::OcrError::Failed(e)) => Err(format!("error: ocr: {e}")),
};
if tx.send((space_id, name, OcrUpdate::Done(result))).is_err() {
return;
}
}
});
}
pub(crate) fn ocr_backend(&self) -> Option<OcrBackend> {
if self.ocr_engine == "local" {
let model = self.local_ocr_model.trim();
let model = if model.is_empty() { "glm-ocr" } else { model };
return Some(OcrBackend::Ollama(
reqwest::Client::new(),
model.to_string(),
));
}
if self.vlm_ocr_enabled() {
let model = self.ocr_model.trim().to_string();
return self
.resolve_model_backend(&model)
.map(|(p, raw_model)| OcrBackend::Router(p, raw_model));
}
None
}
pub(crate) fn ocr_local_install(&mut self, arg: &str) {
if self.ocr_pull_rx.is_some() {
self.status = "an OCR model pull is already running".to_string();
return;
}
let model = if arg.is_empty() {
"glm-ocr".to_string()
} else {
arg.to_string()
};
self.local_ocr_model.clone_from(&model);
let _ = self.db.set_setting("local_ocr_model", &model);
#[cfg(test)]
{
self.ocr_engine = "local".to_string();
let _ = self.db.set_setting("ocr_engine", "local");
self.status = format!("(test) local OCR: {model}");
}
#[cfg(not(test))]
{
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
self.ocr_pull_rx = Some(rx);
self.status = format!("pulling {model} via ollama… (keeps running in background)");
tokio::spawn(async move {
let result = match tokio::process::Command::new("ollama")
.args(["pull", &model])
.output()
.await
{
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
Err("ollama not installed — get it from https://ollama.com (pacman -S ollama), then retry".to_string())
}
Err(e) => Err(format!("ollama pull failed: {e}")),
Ok(out) if !out.status.success() => {
let err = String::from_utf8_lossy(&out.stderr);
let hint = if err.contains("could not connect") || err.contains("connection refused") {
" — is the ollama server running? (systemctl start ollama, or `ollama serve`)"
} else {
""
};
Err(format!("ollama pull failed: {}{hint}", err.trim()))
}
Ok(_) => Ok(model),
};
let _ = tx.send(result);
});
}
}
pub fn on_ocr_pull(&mut self, r: Option<Result<String, String>>) {
let Some(result) = r else {
self.ocr_pull_rx = None;
return;
};
self.ocr_pull_rx = None;
match result {
Ok(model) => {
self.ocr_engine = "local".to_string();
let _ = self.db.set_setting("ocr_engine", "local");
self.status = format!(
"local OCR ready: {model} via ollama — Ctrl+O a file in /files to re-run it"
);
}
Err(e) => self.status = e,
}
}
pub(crate) fn reextract_selected_file(&mut self) {
let Some(f) = self.files_cache.get(self.files_selected).cloned() else {
return;
};
let _ = self.db.set_file_chunks(&f.id, &[]);
let _ = self
.db
.upsert_file(&self.active_space.id, &f.name, "", 0, "re-extracting");
self.status = format!("re-extracting: {}", f.name);
self.rescan_files();
}
pub(crate) fn reocr_selected_file(&mut self) {
let Some(f) = self.files_cache.get(self.files_selected).cloned() else {
return;
};
let ext = std::path::Path::new(&f.name)
.extension()
.and_then(|e| e.to_str())
.unwrap_or("")
.to_lowercase();
if ext != "pdf" && !crate::extract::is_image_ext(&ext) {
self.status = format!("only PDFs and images support OCR: {}", f.name);
return;
}
let path = self.space.files_dir(&self.active_space.name).join(&f.name);
self.ocr_rx = None;
let _ = self.db.set_file_status(&f.id, "ocr…");
self.start_ocr(vec![(self.active_space.id.clone(), f.name.clone(), path)]);
self.files_cache = self
.db
.list_files(&self.active_space.id)
.unwrap_or_default();
self.status = format!("ocr queued: {}", f.name);
}
pub(crate) fn start_embedding(&mut self) {
if self.embed_rx.is_some() {
return;
}
let model = self.embedding_model.trim().to_string();
if model.is_empty() {
return;
}
let Some((provider, raw_model)) = self.resolve_model_backend(&model) else {
return;
};
let space_id = self.active_space.id.clone();
let Ok(missing) = self.db.files_missing_embeddings(&space_id) else {
return;
};
let Some(file_id) = missing.first().cloned() else {
return;
};
let chunks = self.db.file_chunk_texts(&file_id).unwrap_or_default();
if chunks.is_empty() {
return;
}
let Ok(handle) = tokio::runtime::Handle::try_current() else {
return;
};
let _ = self.db.set_file_status(&file_id, "embedding…");
if space_id == self.active_space.id {
self.files_cache = self.db.list_files(&space_id).unwrap_or_default();
}
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
self.embed_rx = Some(rx);
handle.spawn(async move {
let mut out: Vec<(i64, Vec<f32>)> = Vec::with_capacity(chunks.len());
let mut err = None;
for batch in chunks.chunks(64) {
let inputs: Vec<String> = batch.iter().map(|(_, t)| t.clone()).collect();
match provider.embed(&raw_model, inputs).await {
Ok(vecs) => out.extend(batch.iter().zip(vecs).map(|((seq, _), v)| (*seq, v))),
Err(e) => {
err = Some(e.to_string());
break;
}
}
}
let result = match err {
Some(e) => Err(e),
None => Ok(out),
};
let _ = tx.send((space_id, file_id, result));
});
}
pub fn on_embed_done(&mut self, r: Option<crate::app::EmbedMsg>) {
let Some((space_id, file_id, result)) = r else {
self.embed_rx = None;
return;
};
self.embed_rx = None;
match result {
Ok(vecs) => {
let _ = self.db.set_chunk_embeddings(&file_id, &vecs);
let _ = self.db.set_file_status(&file_id, "ok");
self.start_embedding();
}
Err(e) => {
let _ = self.db.set_file_status(&file_id, "ok");
self.status = format!("embedding failed: {e}");
}
}
if space_id == self.active_space.id {
self.files_cache = self.db.list_files(&space_id).unwrap_or_default();
}
}
pub fn on_ocr_done(&mut self, r: Option<(String, String, OcrUpdate)>) {
let Some((space_id, name, update)) = r else {
self.ocr_rx = None;
self.rescan_files();
return;
};
let Ok(files) = self.db.list_files(&space_id) else {
return;
};
let Some(f) = files.iter().find(|f| f.name == name) else {
return; };
if !f.status.starts_with("ocr") {
return; }
match update {
OcrUpdate::Stage(s) => {
let _ = self.db.set_file_status(&f.id, &format!("ocr: {s}"));
if space_id == self.active_space.id {
self.status = format!("ocr {name}: {s}");
}
}
OcrUpdate::Progress(done, total, failed) => {
let tail = if failed > 0 {
format!(" ({failed} failed)")
} else {
String::new()
};
let _ = self
.db
.set_file_status(&f.id, &format!("ocr {done}/{total}{tail}"));
if space_id == self.active_space.id {
self.status = format!("ocr {name}: {done}/{total} pages{tail}");
}
}
OcrUpdate::Done(Ok((text, errors))) if text.trim().is_empty() => {
let status = match errors.first() {
Some((i, e)) => {
format!("all pages failed (p{}: {})", i + 1, clip_err(e))
}
None => "no text (ocr found nothing)".to_string(),
};
let _ = self.db.set_file_status(&f.id, &status);
if space_id == self.active_space.id {
self.status = format!("ocr {name}: {status}");
}
}
OcrUpdate::Done(Ok((text, errors))) => {
let _ = self
.db
.set_file_chunks(&f.id, &crate::extract::chunk_lines(&text));
let status = match errors.first() {
None => "ok".to_string(),
Some((i, e)) => format!(
"ok — {} page{} failed (p{}: {})",
errors.len(),
if errors.len() == 1 { "" } else { "s" },
i + 1,
clip_err(e),
),
};
let _ = self.db.set_file_status(&f.id, &status);
if let Some(new_name) = Self::descriptive_paste_name(f, &text) {
let dir = self.space.files_dir(&self.active_space.name);
let old_path = dir.join(&f.name);
let new_path = dir.join(&new_name);
if old_path.exists() && std::fs::rename(&old_path, &new_path).is_ok() {
let _ = self.db.rename_file(&f.id, &new_name);
let _ = self
.db
.replace_file_ref_in_messages(&space_id, &f.name, &new_name);
if space_id == self.active_space.id {
self.status = format!("ocr done: {new_name}");
}
} else if space_id == self.active_space.id {
self.status = format!("ocr done: {name}");
}
} else if space_id == self.active_space.id {
self.status = format!("ocr done: {name}");
}
}
OcrUpdate::Done(Err(msg)) => {
let _ = self.db.set_file_status(&f.id, &msg);
if space_id == self.active_space.id {
self.status = format!("ocr {name}: {msg}");
}
}
}
if space_id == self.active_space.id {
self.files_cache = self.db.list_files(&space_id).unwrap_or_default();
self.files_selected = self
.files_selected
.min(self.files_cache.len().saturating_sub(1));
}
}
pub fn import_file(&mut self, path: &Path) -> Result<String> {
let name = path
.file_name()
.map(|n| n.to_string_lossy().to_string())
.filter(|n| !n.is_empty())
.context("path has no file name")?;
let dir = self.space.files_dir(&self.active_space.name);
std::fs::create_dir_all(&dir).with_context(|| format!("creating {}", dir.display()))?;
std::fs::copy(path, dir.join(&name))
.with_context(|| format!("copying {} into the space", path.display()))?;
self.rescan_files();
Ok(name)
}
pub fn confirm_files_delete(&mut self) -> Result<()> {
if let Some(f) = self.files_cache.get(self.files_selected).cloned() {
let disk = self.space.files_dir(&self.active_space.name).join(&f.name);
if disk.exists() {
std::fs::remove_file(&disk)
.with_context(|| format!("removing {}", disk.display()))?;
}
self.db.delete_file(&f.id)?;
self.status = format!("removed {}", f.name);
self.rescan_files();
}
self.files_mode = super::FilesMode::Browse;
Ok(())
}
pub(crate) fn open_files_popup(&mut self, tab: super::FilesTab) {
self.files_tab = tab;
match tab {
super::FilesTab::Images => {
self.refresh_images();
}
super::FilesTab::Scripts => {
self.refresh_scripts();
}
super::FilesTab::Files => {
self.rescan_files();
}
}
self.files_mode = super::FilesMode::Browse;
self.popup = super::Popup::Files;
}
pub fn move_files_selection(&mut self, delta: i32) {
self.files_selected =
super::clamp_cursor(self.files_selected, self.files_cache.len(), delta);
}
pub fn start_files_add(&mut self) {
self.files_edit.clear();
self.files_mode = super::FilesMode::Add;
}
pub fn confirm_files_add(&mut self) {
let raw = self.files_edit.trim().to_string();
self.files_mode = super::FilesMode::Browse;
if raw.is_empty() {
return;
}
let path = std::path::PathBuf::from(&raw);
if !path.is_file() {
self.status = format!("not a file: {raw}");
return;
}
match self.import_file(&path) {
Ok(name) => self.status = format!("imported {name}"),
Err(e) => self.status = format!("import failed: {e}"),
}
}
pub fn start_files_rename(&mut self) {
if let Some(f) = self.files_cache.get(self.files_selected) {
self.files_edit = f.name.clone();
self.files_mode = super::FilesMode::Rename;
}
}
pub fn confirm_files_rename(&mut self) -> Result<()> {
let new = self.files_edit.trim().to_string();
self.files_mode = super::FilesMode::Browse;
let Some(f) = self.files_cache.get(self.files_selected).cloned() else {
return Ok(());
};
if new.is_empty() || new == f.name {
return Ok(());
}
if new.contains(['/', '\\']) || new == "." || new == ".." {
self.status = format!("invalid name: {new}");
return Ok(());
}
let dir = self.space.files_dir(&self.active_space.name);
if dir.join(&new).exists() {
self.status = format!("{new} already exists");
return Ok(());
}
std::fs::rename(dir.join(&f.name), dir.join(&new))
.with_context(|| format!("renaming {} to {new}", f.name))?;
self.rescan_files();
self.files_selected = self
.files_cache
.iter()
.position(|f| f.name == new)
.unwrap_or(self.files_selected);
self.status = format!("renamed {} to {new}", f.name);
Ok(())
}
pub fn open_selected_file(&mut self) {
if let Some(f) = self.files_cache.get(self.files_selected) {
let path = self.space.files_dir(&self.active_space.name).join(&f.name);
let ext = path.extension().and_then(|e| e.to_str()).unwrap_or("");
if ext == "md" {
self.pending_editor = Some(super::PendingEditor::ScriptFile(path));
} else {
let _ = open::that_detached(&path);
}
self.status = format!("opened {}", f.name);
}
}
fn descriptive_paste_name(f: &FileRow, ocr_text: &str) -> Option<String> {
let stem = std::path::Path::new(&f.name).file_stem()?.to_str()?;
let ext = std::path::Path::new(&f.name).extension()?.to_str()?;
if !is_uuid_like(stem) {
return None;
}
let slug = Self::slug_from_ocr(ocr_text)?;
Some(format!("{stem}-{slug}.{ext}"))
}
fn slug_from_ocr(text: &str) -> Option<String> {
let words: Vec<&str> = text
.split_whitespace()
.filter(|w| {
let w = w.trim_matches(|c: char| !c.is_alphanumeric());
w.len() > 2 && w.chars().any(char::is_alphanumeric)
})
.collect();
if words.is_empty() {
return None;
}
let slug: String = words
.iter()
.take(5)
.map(|w| {
w.trim_matches(|c: char| !c.is_alphanumeric())
.to_lowercase()
})
.collect::<Vec<_>>()
.join("_");
if slug.is_empty() { None } else { Some(slug) }
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::db::Db;
use crate::space::Space;
fn test_app() -> App {
let db = Db::open_in_memory().unwrap();
let root = std::env::temp_dir().join(format!("nexus-files-test-{}", uuid::Uuid::new_v4()));
std::fs::create_dir_all(root.join("spaces")).unwrap();
let space = Space { root };
App::new(db, Some("k"), space)
}
#[tokio::test]
async fn embedder_queue_backfills_chains_and_stops_on_error() {
let mut a = test_app();
let space = a.active_space.id.clone();
let id = a.db.upsert_file(&space, "b.txt", "h", 1, "ok").unwrap();
a.db.set_file_chunks(&id, &[("l".into(), "text".into())])
.unwrap();
let saved = a.backends.clone();
a.backends = crate::app::Backends::default();
a.start_embedding();
assert!(a.embed_rx.is_none());
a.backends = saved;
let m = std::mem::take(&mut a.embedding_model);
a.start_embedding();
assert!(a.embed_rx.is_none());
a.embedding_model = m;
a.start_embedding();
assert!(a.embed_rx.is_some());
let files = a.db.list_files(&space).unwrap();
assert!(
files[0].status.starts_with("embedding"),
"{}",
files[0].status
);
a.on_embed_done(Some((
space.clone(),
id.clone(),
Ok(vec![(0, vec![1.0f32, 0.0])]),
)));
assert!(a.db.files_missing_embeddings(&space).unwrap().is_empty());
assert_eq!(a.db.list_files(&space).unwrap()[0].status, "ok");
a.db.set_file_chunks(&id, &[("l".into(), "new".into())])
.unwrap();
a.on_embed_done(Some((space.clone(), id.clone(), Err("offline".into()))));
assert!(a.embed_rx.is_none());
assert!(a.status.contains("embedding failed"));
assert_eq!(a.db.list_files(&space).unwrap()[0].status, "ok");
}
#[test]
fn import_copies_extracts_and_indexes() {
let mut a = test_app();
let src = std::env::temp_dir().join(format!("nexus-src-{}.md", uuid::Uuid::new_v4()));
std::fs::write(&src, "# quarterly report\nrevenue up").unwrap();
let name = a.import_file(&src).unwrap();
assert_eq!(name, src.file_name().unwrap().to_string_lossy());
assert_eq!(a.files_cache.len(), 1);
assert_eq!(a.files_cache[0].status, "ok");
assert!(a.space.files_dir(&a.active_space.name).join(&name).exists());
let hits = crate::db::search_chunks(a.db.conn_for_test(), &a.active_space.id, "revenue", 8)
.unwrap();
assert_eq!(hits.len(), 1);
}
#[test]
fn ollama_ocr_body_uses_native_generate_shape() {
let body = ollama_ocr_body("glm-ocr", "QUFB");
assert_eq!(body["model"], "glm-ocr");
assert_eq!(body["stream"], false);
assert_eq!(body["images"][0], "QUFB"); assert!(body["prompt"].as_str().unwrap().contains("furigana"));
assert!(
body.get("messages").is_none(),
"must not be OpenAI chat shape"
);
}
#[test]
fn ocr_backend_routes_by_engine() {
let mut a = test_app();
assert!(matches!(a.ocr_backend(), Some(OcrBackend::Router(..))));
a.ocr_engine = "local".to_string();
a.local_ocr_model = String::new(); match a.ocr_backend() {
Some(OcrBackend::Ollama(_, model)) => assert_eq!(model, "glm-ocr"),
other => panic!("expected ollama backend, got {}", other.is_some()),
}
a.ocr_engine = "tesseract".to_string();
assert!(a.ocr_backend().is_none());
a.ocr_engine = "auto".to_string();
a.backends = crate::app::Backends::default();
assert!(a.ocr_backend().is_none());
}
#[tokio::test]
async fn ocr_pull_success_switches_engine_to_local() {
let mut a = test_app();
a.on_ocr_pull(Some(Ok("glm-ocr".to_string())));
assert_eq!(a.ocr_engine, "local");
assert!(a.status.contains("local OCR ready"));
a.on_ocr_pull(Some(Err("ollama not installed — get it".to_string())));
assert_eq!(a.ocr_engine, "local"); assert!(a.status.contains("ollama not installed"));
}
#[test]
fn ocr_statuses_surface_stages_failures_and_reasons() {
let mut a = test_app();
let space = a.active_space.id.clone();
let id =
a.db.upsert_file(&space, "scan.pdf", "h", 1, "ocr…")
.unwrap();
a.on_ocr_done(Some((
space.clone(),
"scan.pdf".into(),
OcrUpdate::Stage("rendering pages (300 dpi)…".into()),
)));
let status = a.db.list_files(&space).unwrap()[0].status.clone();
assert_eq!(status, "ocr: rendering pages (300 dpi)…");
a.on_ocr_done(Some((
space.clone(),
"scan.pdf".into(),
OcrUpdate::Progress(5, 10, 2),
)));
assert_eq!(
a.db.list_files(&space).unwrap()[0].status,
"ocr 5/10 (2 failed)"
);
assert!(a.status.contains("5/10 pages (2 failed)"));
a.on_ocr_done(Some((
space.clone(),
"scan.pdf".into(),
OcrUpdate::Done(Ok((
"[page 1]\ntext".to_string(),
vec![
(2, "timeout after 600s".to_string()),
(4, "boom".to_string()),
],
))),
)));
let status = a.db.list_files(&space).unwrap()[0].status.clone();
assert_eq!(status, "ok — 2 pages failed (p3: timeout after 600s)");
let _ = a.db.set_file_status(&id, "ocr…");
a.on_ocr_done(Some((
space.clone(),
"scan.pdf".into(),
OcrUpdate::Done(Ok((
String::new(),
vec![(0, "cannot reach ollama at 127.0.0.1:11434 — is it running? (systemctl start ollama)".to_string())],
))),
)));
let status = a.db.list_files(&space).unwrap()[0].status.clone();
assert!(
status.starts_with("all pages failed (p1: cannot reach ollama"),
"{status}"
);
assert!(!status.starts_with("ocr"), "{status}");
}
#[test]
fn reextract_clears_stale_chunks_and_reindexes_from_disk() {
let mut a = test_app();
let dir = a.space.files_dir(&a.active_space.name);
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("doc.txt"), "real content on disk").unwrap();
a.rescan_files();
let id = a.files_cache[0].id.clone();
a.db.set_file_chunks(&id, &[("p1".into(), "garbage".into())])
.unwrap();
a.files_selected = 0;
a.reextract_selected_file();
assert!(a.status.contains("re-extracting"), "{}", a.status);
let texts = a.db.file_chunk_texts(&id).unwrap();
assert_eq!(texts.len(), 1);
assert!(texts[0].1.contains("real content"), "{texts:?}");
assert_eq!(a.files_cache[0].status, "ok");
}
#[test]
fn rescan_picks_up_dropped_and_deleted_files() {
let mut a = test_app();
let dir = a.space.files_dir(&a.active_space.name);
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("dropped.txt"), "hello dropped").unwrap();
a.rescan_files();
assert_eq!(a.files_cache.len(), 1);
assert_eq!(a.files_cache[0].name, "dropped.txt");
std::fs::write(dir.join("dropped.txt"), "hello again").unwrap();
a.rescan_files();
assert_eq!(a.files_cache.len(), 1);
std::fs::remove_file(dir.join("dropped.txt")).unwrap();
a.rescan_files();
assert!(a.files_cache.is_empty());
}
#[test]
fn empty_extraction_gets_no_text_status() {
let mut a = test_app();
let dir = a.space.files_dir(&a.active_space.name);
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("empty.txt"), " ").unwrap();
a.rescan_files();
assert_eq!(a.files_cache[0].status, "no text (scanned?)");
}
#[test]
fn rescan_skips_stat_unchanged_files_without_rehashing() {
let mut a = test_app();
let dir = a.space.files_dir(&a.active_space.name);
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("book.txt"), "big content").unwrap();
a.rescan_files();
let f = a.files_cache[0].clone();
assert!(f.mtime > 0, "mtime recorded on index");
a.db.upsert_file(&a.active_space.id, "book.txt", "sentinel", f.size, "ok")
.unwrap();
a.rescan_files();
assert_eq!(a.files_cache[0].hash, "sentinel");
std::fs::write(dir.join("book.txt"), "big content grew").unwrap();
a.rescan_files();
assert_ne!(a.files_cache[0].hash, "sentinel");
assert_eq!(a.files_cache[0].status, "ok");
}
#[test]
fn rename_moves_disk_file_and_reindexes() {
let mut a = test_app();
let dir = a.space.files_dir(&a.active_space.name);
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("old.txt"), "searchable content").unwrap();
std::fs::write(dir.join("taken.txt"), "x").unwrap();
a.rescan_files();
a.files_selected = a
.files_cache
.iter()
.position(|f| f.name == "old.txt")
.unwrap();
a.start_files_rename();
assert_eq!(a.files_edit, "old.txt");
a.files_edit = "taken.txt".to_string();
a.confirm_files_rename().unwrap();
assert!(a.status.contains("already exists"));
assert!(dir.join("old.txt").exists());
a.files_selected = a
.files_cache
.iter()
.position(|f| f.name == "old.txt")
.unwrap();
a.start_files_rename();
a.files_edit = "../evil.txt".to_string();
a.confirm_files_rename().unwrap();
assert!(a.status.contains("invalid name"));
a.files_selected = a
.files_cache
.iter()
.position(|f| f.name == "old.txt")
.unwrap();
a.start_files_rename();
a.files_edit = "new.txt".to_string();
a.confirm_files_rename().unwrap();
assert!(!dir.join("old.txt").exists());
assert!(dir.join("new.txt").exists());
assert!(a.files_cache.iter().any(|f| f.name == "new.txt"));
assert_eq!(a.files_cache[a.files_selected].name, "new.txt");
}
#[test]
fn delete_removes_disk_file_and_row() {
let mut a = test_app();
let dir = a.space.files_dir(&a.active_space.name);
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("gone.txt"), "bye").unwrap();
a.rescan_files();
a.files_selected = 0;
a.confirm_files_delete().unwrap();
assert!(a.files_cache.is_empty());
assert!(!dir.join("gone.txt").exists());
}
#[test]
fn files_command_opens_popup_and_rescans() {
let mut a = test_app();
let dir = a.space.files_dir(&a.active_space.name);
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("seen.txt"), "content").unwrap();
a.run_command("files").unwrap();
assert!(a.popup == crate::app::Popup::Files);
assert_eq!(a.files_cache.len(), 1);
assert!(a.files_mode == crate::app::FilesMode::Browse);
}
#[test]
fn confirm_files_add_imports_typed_path() {
let mut a = test_app();
let src = std::env::temp_dir().join(format!("nexus-add-{}.txt", uuid::Uuid::new_v4()));
std::fs::write(&src, "typed in").unwrap();
a.start_files_add();
assert!(a.files_mode == crate::app::FilesMode::Add);
a.files_edit = src.to_string_lossy().to_string();
a.confirm_files_add();
assert!(a.files_mode == crate::app::FilesMode::Browse);
assert_eq!(a.files_cache.len(), 1);
a.start_files_add();
a.files_edit = "/definitely/not/a/file".to_string();
a.confirm_files_add();
assert!(a.status.contains("not a file"));
assert_eq!(a.files_cache.len(), 1);
}
#[test]
fn picker_lists_dirs_first_descends_and_imports() {
let mut a = test_app();
let root = std::env::temp_dir().join(format!("nexus-pick-{}", uuid::Uuid::new_v4()));
std::fs::create_dir_all(root.join("subdir")).unwrap();
std::fs::write(root.join("bbb.txt"), "file b").unwrap();
std::fs::write(root.join("aaa.txt"), "file a").unwrap();
a.picker_dir = root.clone();
a.open_file_picker();
assert!(a.files_mode == crate::app::FilesMode::Pick);
let names: Vec<&str> = a
.filtered_picker_entries()
.iter()
.map(|e| e.name.as_str())
.collect();
assert_eq!(names, vec!["subdir", "aaa.txt", "bbb.txt"]);
a.picker_selected = 0;
a.picker_enter();
assert_eq!(a.picker_dir, root.join("subdir"));
assert!(a.filtered_picker_entries().is_empty());
a.picker_backspace();
assert_eq!(a.picker_dir, root);
let idx = a
.filtered_picker_entries()
.iter()
.position(|e| e.name == "aaa.txt")
.unwrap();
a.picker_selected = idx;
a.picker_enter();
assert!(a.files_mode == crate::app::FilesMode::Browse);
assert!(a.files_cache.iter().any(|f| f.name == "aaa.txt"));
}
#[test]
fn picker_filter_fuzzy_matches_and_backspace_edits_filter_first() {
let mut a = test_app();
let root = std::env::temp_dir().join(format!("nexus-pick-{}", uuid::Uuid::new_v4()));
std::fs::create_dir_all(&root).unwrap();
std::fs::write(root.join("report-2026.pdf"), "x").unwrap();
std::fs::write(root.join("notes.md"), "y").unwrap();
a.picker_dir = root.clone();
a.open_file_picker();
a.picker_filter_push('r');
a.picker_filter_push('p');
a.picker_filter_push('t');
let names: Vec<&str> = a
.filtered_picker_entries()
.iter()
.map(|e| e.name.as_str())
.collect();
assert_eq!(names, vec!["report-2026.pdf"]);
a.picker_backspace();
assert_eq!(a.picker_filter, "rp");
assert_eq!(a.picker_dir, root);
}
#[tokio::test]
async fn rescan_marks_empty_pdf_ocr_and_spawns_batch() {
let mut a = test_app();
let dir = a.space.files_dir(&a.active_space.name);
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("scan.pdf"), crate::extract::minimal_pdf(None)).unwrap();
a.rescan_files();
assert_eq!(a.files_cache[0].status, "ocr…");
assert!(a.ocr_rx.is_some(), "an ocr batch should be in flight");
a.rescan_files();
assert_eq!(a.files_cache[0].status, "ocr…");
}
#[tokio::test]
async fn rescan_requeues_stale_ocr_status_when_idle() {
let mut a = test_app();
let dir = a.space.files_dir(&a.active_space.name);
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("scan.pdf"), crate::extract::minimal_pdf(None)).unwrap();
a.rescan_files();
a.ocr_rx = None;
a.rescan_files();
assert!(a.ocr_rx.is_some(), "stale ocr… should re-queue");
}
#[test]
fn on_ocr_done_ok_indexes_and_marks_ok() {
let mut a = test_app();
let id =
a.db.upsert_file(&a.active_space.id, "scan.pdf", "h", 9, "ocr…")
.unwrap();
let _ = id;
a.on_ocr_done(Some((
a.active_space.id.clone(),
"scan.pdf".to_string(),
OcrUpdate::Done(Ok((
"[page 1]\nquarterly revenue table".to_string(),
Vec::new(),
))),
)));
assert_eq!(a.files_cache[0].status, "ok");
let hits = crate::db::search_chunks(a.db.conn_for_test(), &a.active_space.id, "revenue", 8)
.unwrap();
assert_eq!(hits.len(), 1);
}
#[test]
fn on_ocr_progress_updates_status_and_status_line() {
let mut a = test_app();
a.db.upsert_file(&a.active_space.id, "scan.pdf", "h", 9, "ocr…")
.unwrap();
a.on_ocr_done(Some((
a.active_space.id.clone(),
"scan.pdf".to_string(),
OcrUpdate::Progress(3, 10, 0),
)));
assert_eq!(a.files_cache[0].status, "ocr 3/10");
assert!(a.status.contains("3/10"), "{}", a.status);
a.on_ocr_done(Some((
a.active_space.id.clone(),
"scan.pdf".to_string(),
OcrUpdate::Done(Ok(("[page 1]\nfound".to_string(), Vec::new()))),
)));
assert_eq!(a.files_cache[0].status, "ok");
}
#[test]
fn on_ocr_done_empty_and_err_statuses() {
let mut a = test_app();
a.db.upsert_file(&a.active_space.id, "blank.pdf", "h1", 9, "ocr…")
.unwrap();
a.db.upsert_file(&a.active_space.id, "bad.pdf", "h2", 9, "ocr…")
.unwrap();
a.on_ocr_done(Some((
a.active_space.id.clone(),
"blank.pdf".to_string(),
OcrUpdate::Done(Ok((String::new(), Vec::new()))),
)));
a.on_ocr_done(Some((
a.active_space.id.clone(),
"bad.pdf".to_string(),
OcrUpdate::Done(Err(
"scanned pdf — install tesseract + poppler for ocr".to_string()
)),
)));
let by_name = |a: &App, n: &str| {
a.files_cache
.iter()
.find(|f| f.name == n)
.unwrap()
.status
.clone()
};
assert_eq!(by_name(&a, "blank.pdf"), "no text (ocr found nothing)");
assert_eq!(
by_name(&a, "bad.pdf"),
"scanned pdf — install tesseract + poppler for ocr"
);
}
#[test]
fn on_ocr_done_for_inactive_space_writes_db_but_not_cache() {
let mut a = test_app();
let other = a.db.create_space("other").unwrap();
a.db.upsert_file(&other.id, "scan.pdf", "h", 9, "ocr…")
.unwrap();
a.on_ocr_done(Some((
other.id.clone(),
"scan.pdf".to_string(),
OcrUpdate::Done(Ok(("found text".to_string(), Vec::new()))),
)));
assert!(
a.files_cache.is_empty(),
"active-space cache must not show other space's file"
);
let rows = a.db.list_files(&other.id).unwrap();
assert_eq!(rows[0].status, "ok");
a.on_ocr_done(Some((
other.id.clone(),
"gone.pdf".to_string(),
OcrUpdate::Done(Ok(("x".to_string(), Vec::new()))),
)));
}
#[tokio::test]
async fn on_ocr_done_none_clears_channel_and_requeues_stragglers() {
let mut a = test_app();
let dir = a.space.files_dir(&a.active_space.name);
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("scan.pdf"), crate::extract::minimal_pdf(None)).unwrap();
a.rescan_files();
assert!(a.ocr_rx.is_some());
a.on_ocr_done(None);
assert!(a.ocr_rx.is_some(), "straggler should re-queue on batch end");
}
}