use super::*;
use std::fs;
use std::path::PathBuf;
fn scratch_dir(tag: &str) -> PathBuf {
let pid = std::process::id();
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0);
let p = std::env::temp_dir().join(format!("trusty-search-index-{tag}-{pid}-{nanos}"));
let _ = fs::remove_dir_all(&p);
p
}
#[test]
fn ensure_project_indexed_withholds_id_when_nothing_was_registered() {
let _guard = crate::data_dir::ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let data_dir = scratch_dir("data");
fs::create_dir_all(&data_dir).unwrap();
unsafe {
std::env::set_var(crate::data_dir::DATA_DIR_OVERRIDE_ENV, &data_dir);
}
let project = scratch_dir("git");
fs::create_dir_all(project.join(".git")).unwrap();
let nested = project.join("crates/inner");
fs::create_dir_all(&nested).unwrap();
let report = ensure_project_indexed_reporting(
&nested,
IndexOptions::default().with_allow_sensitive_path(true),
);
let pinnable = ensure_project_indexed(&nested, true);
let expected = crate::derive_index_id(&project);
unsafe {
std::env::remove_var(crate::data_dir::DATA_DIR_OVERRIDE_ENV);
}
let _ = fs::remove_dir_all(&project);
let _ = fs::remove_dir_all(&data_dir);
assert_eq!(
report.index_id,
Some(expected),
"id is the git-root basename"
);
assert_ne!(
report.registration,
IndexRegistration::Confirmed,
"no daemon was contacted, so nothing can be confirmed"
);
assert_eq!(
pinnable, None,
"an unregistered index must not come back as a pinnable id (#5091)"
);
}
#[test]
fn reporting_says_skipped_under_test_harness() {
let _guard = crate::data_dir::ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let project = scratch_dir("report-skip");
fs::create_dir_all(project.join(".git")).unwrap();
let report = ensure_project_indexed_reporting(&project, IndexOptions::default());
let _ = fs::remove_dir_all(&project);
assert_eq!(
report.registration,
IndexRegistration::SkippedUnderTest,
"a test process suppresses the write (#4255) and must say so"
);
assert!(
report.index_id.is_some(),
"the id is still returned — the fail-open contract is unchanged"
);
}
#[test]
fn reporting_says_daemon_unreachable_when_no_daemon_is_discoverable() {
let _guard = crate::data_dir::ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let data_dir = scratch_dir("report-nodaemon-data");
fs::create_dir_all(&data_dir).unwrap();
let project = scratch_dir("report-nodaemon");
fs::create_dir_all(project.join(".git")).unwrap();
unsafe {
std::env::set_var(crate::data_dir::DATA_DIR_OVERRIDE_ENV, &data_dir);
std::env::set_var(crate::test_harness::ALLOW_PRODUCTION_ENV, "1");
}
let report = ensure_project_indexed_reporting(&project, IndexOptions::default());
unsafe {
std::env::remove_var(crate::test_harness::ALLOW_PRODUCTION_ENV);
std::env::remove_var(crate::data_dir::DATA_DIR_OVERRIDE_ENV);
}
let _ = fs::remove_dir_all(&project);
let _ = fs::remove_dir_all(&data_dir);
assert_eq!(
report.registration,
IndexRegistration::DaemonUnreachable,
"no address file means nothing was sent — that is not a registration"
);
assert!(report.index_id.is_some(), "the id is still returned");
}
#[test]
fn ensure_project_indexed_none_for_root() {
assert_eq!(ensure_project_indexed(Path::new("/"), true), None);
assert_eq!(ensure_project_indexed(Path::new("/"), false), None);
assert_eq!(
ensure_project_indexed_reporting(Path::new("/"), IndexOptions::default()).registration,
IndexRegistration::RefusedUnindexableRoot(crate::IndexRootRefusal::FilesystemRoot)
);
}
#[test]
fn ensure_project_indexed_refuses_the_real_home_directory() {
let Some(home) = dirs::home_dir() else {
panic!("this test needs a resolvable home directory");
};
assert_eq!(
crate::resolve_project_root(&home),
home,
"no ancestor of $HOME may be a git repository for this case to exist"
);
let report = ensure_project_indexed_reporting(&home, IndexOptions::default());
assert_eq!(
report.registration,
IndexRegistration::RefusedUnindexableRoot(crate::IndexRootRefusal::HomeDirectory)
);
assert_eq!(
report.index_id, None,
"the home directory's basename is the wrong id and must not be handed back"
);
assert_eq!(ensure_project_indexed(&home, false), None);
assert_eq!(ensure_project_indexed(&home, true), None);
}
#[test]
fn index_files_inner_refuses_the_real_home_directory() {
let Some(home) = dirs::home_dir() else {
panic!("this test needs a resolvable home directory");
};
index_files_inner(&home, &[PathBuf::from("some/file.rs")]);
}
#[test]
fn index_files_inner_is_noop_for_empty_paths() {
index_files_inner(Path::new("/"), &[]);
}
#[test]
fn index_files_inner_skips_when_index_id_empty() {
index_files_inner(Path::new("/"), &[PathBuf::from("some/file.rs")]);
}
#[test]
fn index_files_inner_skips_gracefully_when_daemon_down() {
let _guard = crate::data_dir::ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let data_dir = scratch_dir("data-incr");
fs::create_dir_all(&data_dir).unwrap();
unsafe {
std::env::set_var(crate::data_dir::DATA_DIR_OVERRIDE_ENV, &data_dir);
}
let project = scratch_dir("git-incr");
fs::create_dir_all(project.join(".git")).unwrap();
fs::write(project.join("main.rs"), "fn main() {}\n").unwrap();
index_files_inner(&project, &[PathBuf::from("main.rs")]);
unsafe {
std::env::remove_var(crate::data_dir::DATA_DIR_OVERRIDE_ENV);
}
let _ = fs::remove_dir_all(&project);
let _ = fs::remove_dir_all(&data_dir);
}
#[test]
fn index_files_best_effort_drops_the_batch_when_the_shared_pool_is_saturated() {
use crate::index_dispatch::{INDEX_QUEUE_CAPACITY, MAX_INDEX_WORKERS, global};
use std::sync::mpsc::channel;
use std::time::Duration;
let wait = Duration::from_secs(30);
let (started_tx, started_rx) = channel();
let mut releases = Vec::with_capacity(MAX_INDEX_WORKERS);
for _ in 0..MAX_INDEX_WORKERS {
let (release_tx, release_rx) = channel::<()>();
releases.push(release_tx);
let started = started_tx.clone();
assert!(
global().try_submit(Box::new(move || {
let _ = started.send(());
let _ = release_rx.recv_timeout(wait);
})),
"the shared pool refused a job before every worker was even busy"
);
}
for i in 0..MAX_INDEX_WORKERS {
started_rx
.recv_timeout(wait)
.unwrap_or_else(|e| panic!("blocker {i} never started: {e}"));
}
let mut filled = 0usize;
while global().try_submit(Box::new(|| {})) {
filled += 1;
assert!(
filled <= INDEX_QUEUE_CAPACITY,
"the queue accepted {filled} jobs, more than its {INDEX_QUEUE_CAPACITY}-slot capacity"
);
}
let before = global().rejected();
index_files_best_effort(Path::new("/nonexistent-2798"), &[PathBuf::from("main.rs")]);
let after = global().rejected();
let stats = index_drop_stats();
for release in &releases {
let _ = release.send(());
}
assert_eq!(
after,
before + 1,
"a batch submitted to a saturated pool must be dropped and counted"
);
assert_eq!(
stats.dropped_batches, after,
"the public stats must read the same counter the pool increments"
);
assert!(
stats
.seconds_since_last_drop
.is_some_and(|since| since <= 60),
"a drop that just happened must be reported as recent, got {:?}",
stats.seconds_since_last_drop
);
}
#[test]
fn batch_budget_is_exhausted_at_and_past_the_cap() {
use std::time::Duration;
assert!(!batch_budget_exhausted(Duration::from_secs(0)));
assert!(!batch_budget_exhausted(
BATCH_INDEX_BUDGET - Duration::from_millis(1)
));
assert!(batch_budget_exhausted(BATCH_INDEX_BUDGET));
assert!(batch_budget_exhausted(
BATCH_INDEX_BUDGET + Duration::from_secs(600)
));
}
#[test]
fn a_truncated_batch_is_counted_separately_from_a_dropped_one() {
let before = index_drop_stats().truncated_batches;
assert!(
stop_batch_for_budget(
crate::index_dispatch::global(),
BATCH_INDEX_BUDGET,
"idx",
3,
10
),
"a batch that has spent its budget must be stopped"
);
let after = index_drop_stats();
assert_eq!(
after.truncated_batches,
before + 1,
"stopping on the budget must be counted, not only logged"
);
assert!(
after
.seconds_since_last_truncation
.is_some_and(|since| since <= 60),
"a truncation that just happened must be reported as recent, got {:?}",
after.seconds_since_last_truncation
);
}
#[test]
fn an_unexhausted_budget_records_no_truncation() {
use crate::index_dispatch::BoundedDispatcher;
use std::time::Duration;
let pool = BoundedDispatcher::new(1, 1);
assert!(
!stop_batch_for_budget(&pool, Duration::from_secs(0), "idx", 0, 10),
"a batch that has spent none of its budget must not be stopped"
);
assert!(
!stop_batch_for_budget(
&pool,
BATCH_INDEX_BUDGET - Duration::from_millis(1),
"idx",
9,
10
),
"a batch one millisecond inside its budget must not be stopped"
);
assert_eq!(
pool.truncated(),
0,
"a batch that was never stopped must not be counted as truncated"
);
assert_eq!(
pool.last_truncation_unix_secs(),
None,
"with no truncation the stamp must stay unset, never a misleading epoch"
);
}
#[test]
fn relative_index_path_strips_root_prefix() {
let root = Path::new("/Users/dev/my-project");
let abs = root.join("src/main.rs");
assert_eq!(relative_index_path(root, &abs), "src/main.rs");
}
#[test]
fn relative_index_path_falls_back_for_paths_outside_root() {
let root = Path::new("/Users/dev/my-project");
let elsewhere = Path::new("/somewhere/else/file.py");
assert_eq!(
relative_index_path(root, elsewhere),
"/somewhere/else/file.py"
);
}
#[test]
fn index_file_request_body_targets_relative_path_and_content() {
let body = index_file_request_body("src/main.rs", "fn main() {}\n");
assert_eq!(
body.get("path").and_then(serde_json::Value::as_str),
Some("src/main.rs")
);
assert_eq!(
body.get("content").and_then(serde_json::Value::as_str),
Some("fn main() {}\n")
);
assert!(
body.get("allow_sensitive_path").is_none(),
"the per-file endpoint does not re-check the denylist, so no bypass \
flag should be sent: {body:?}"
);
}
#[test]
fn create_index_request_body_respects_allow_sensitive_path_param() {
for root in [
Path::new("/Users/dev/projects/my-repo"),
Path::new("/private/var/folders/xx/scratch-project"),
] {
for allow in [true, false] {
let body = create_index_request_body(
"my-index",
root,
IndexOptions {
allow_sensitive_path: allow,
..IndexOptions::default()
},
);
assert_eq!(
body.get("allow_sensitive_path"),
Some(&serde_json::Value::Bool(allow)),
"request body for root {root:?} must set allow_sensitive_path: {allow}"
);
assert_eq!(
body.get("id").and_then(serde_json::Value::as_str),
Some("my-index")
);
}
}
}
#[test]
fn create_index_request_body_sets_skip_vector() {
let root = Path::new("/Users/dev/projects/my-repo/.worktrees/feat-x");
for allow in [true, false] {
for skip_vector in [true, false] {
let body = create_index_request_body(
"feat-x",
root,
IndexOptions {
allow_sensitive_path: allow,
skip_vector,
},
);
assert_eq!(
body.get("skip_vector"),
Some(&serde_json::Value::Bool(skip_vector)),
"body must set skip_vector: {skip_vector} (allow={allow})"
);
assert_eq!(
body.get("allow_sensitive_path"),
Some(&serde_json::Value::Bool(allow)),
"skip_vector must not disturb allow_sensitive_path"
);
}
}
}
#[test]
fn index_options_default_matches_legacy_ensure_call() {
let root = Path::new("/Users/dev/projects/my-repo");
assert_eq!(
create_index_request_body("my-repo", root, IndexOptions::default()),
create_index_request_body(
"my-repo",
root,
IndexOptions {
allow_sensitive_path: false,
skip_vector: false,
}
)
);
}
#[test]
fn index_is_fresh_true_when_recently_indexed_with_chunks() {
let now = chrono::Utc::now();
let status = serde_json::json!({
"chunk_count": 42,
"last_indexed": now.to_rfc3339(),
});
assert!(index_is_fresh(&status));
}
#[test]
fn index_is_fresh_false_when_no_chunks() {
let now = chrono::Utc::now();
let status = serde_json::json!({
"chunk_count": 0,
"last_indexed": now.to_rfc3339(),
});
assert!(!index_is_fresh(&status));
}
#[test]
fn index_is_fresh_false_when_stale() {
let stale = chrono::Utc::now() - chrono::Duration::hours(2);
let status = serde_json::json!({
"chunk_count": 10,
"last_indexed": stale.to_rfc3339(),
});
assert!(!index_is_fresh(&status));
}
#[test]
fn index_is_fresh_false_when_last_indexed_missing_or_malformed() {
assert!(!index_is_fresh(&serde_json::json!({ "chunk_count": 10 })));
assert!(!index_is_fresh(&serde_json::json!({
"chunk_count": 10,
"last_indexed": "not-a-timestamp",
})));
assert!(!index_is_fresh(&serde_json::json!({})));
}
#[test]
fn retry_backoff_is_bounded_and_increasing() {
use std::time::Duration;
assert_eq!(retry_backoff(1), Duration::from_millis(50));
assert_eq!(retry_backoff(2), Duration::from_millis(150));
assert_eq!(retry_backoff(3), Duration::from_millis(450));
assert!(retry_backoff(2) > retry_backoff(1));
assert!(retry_backoff(3) > retry_backoff(2));
assert_eq!(retry_backoff(100), Duration::from_millis(1000));
}
fn drive_retry_test(
server_fn: impl FnOnce(std::net::TcpListener, std::sync::mpsc::Sender<usize>) + Send + 'static,
) -> (IndexOutcome, usize) {
use std::net::TcpListener;
use std::sync::mpsc;
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let addr = listener.local_addr().unwrap();
let (tx, rx) = mpsc::channel();
let server = std::thread::spawn(move || server_fn(listener, tx));
let client = build_index_client().unwrap();
let url = format!("http://{addr}/indexes/test-index/index-file");
let body = index_file_request_body("src/main.rs", "fn main() {}\n");
let outcome = post_index_file_with_retries(&client, &url, &body);
let accepted = rx
.recv_timeout(std::time::Duration::from_secs(5))
.expect("server thread should have reported an accepted-connection count");
let _ = server.join();
(outcome, accepted)
}
#[test]
fn post_index_file_retries_transient_send_failure() {
use std::io::{Read, Write};
let (outcome, accepted) = drive_retry_test(|listener, tx| {
let mut accepted = 0usize;
for stream in listener.incoming() {
let Ok(mut stream) = stream else { break };
accepted += 1;
if accepted == 1 {
drop(stream);
continue;
}
let mut buf = [0u8; 4096];
let _ = stream.read(&mut buf);
let _ = stream
.write_all(b"HTTP/1.1 200 OK\r\ncontent-length: 0\r\nconnection: close\r\n\r\n");
let _ = stream.flush();
let _ = tx.send(accepted);
break;
}
});
assert_eq!(outcome, IndexOutcome::Indexed);
assert_eq!(accepted, 2);
}
#[test]
fn post_index_file_exhausts_retries_and_returns_send_failed() {
let (outcome, accepted) = drive_retry_test(|listener, tx| {
let mut accepted = 0usize;
for stream in listener.incoming() {
let Ok(stream) = stream else { break };
accepted += 1;
drop(stream);
if accepted >= MAX_INDEX_ATTEMPTS as usize {
let _ = tx.send(accepted);
break;
}
}
});
assert_eq!(outcome, IndexOutcome::SendFailed);
assert_eq!(accepted, MAX_INDEX_ATTEMPTS as usize);
}
fn daemon_was_contacted_during(body: impl FnOnce()) -> bool {
use crate::data_dir::{DATA_DIR_OVERRIDE_ENV, ENV_LOCK};
use std::sync::atomic::{AtomicUsize, Ordering};
let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("bind stand-in daemon");
let addr = listener.local_addr().expect("stand-in daemon local_addr");
listener
.set_nonblocking(true)
.expect("stand-in daemon set_nonblocking");
let guard = ENV_LOCK.lock().unwrap_or_else(|e| e.into_inner());
let data_dir = scratch_dir("4255-daemon");
fs::create_dir_all(&data_dir).expect("create isolated data dir");
let socket_dir = tempfile::tempdir().expect("tempdir for the stand-in socket");
let socket = socket_dir.path().join("s.sock");
let previous = std::env::var(DATA_DIR_OVERRIDE_ENV).ok();
unsafe {
std::env::set_var(DATA_DIR_OVERRIDE_ENV, &data_dir);
std::env::set_var(crate::search_rpc::TRUSTY_SEARCH_SOCKET_ENV, &socket);
}
crate::write_daemon_addr("trusty-search", &addr.to_string()).expect("publish daemon addr");
assert_eq!(
crate::resolve_daemon_base_url("trusty-search"),
Some(format!("http://{addr}")),
"the stand-in HTTP daemon must be discoverable, or this test proves nothing"
);
let calls = std::sync::Arc::new(AtomicUsize::new(0));
let counter = std::sync::Arc::clone(&calls);
let daemon = uds_mock::spawn_blocking_at(socket, move |_method, _params| {
counter.fetch_add(1, Ordering::SeqCst);
Box::pin(async move { Ok(serde_json::json!({ "id": "x", "created": true })) })
});
let derived = crate::search_rpc::search_socket().ok();
assert_eq!(
derived.as_deref(),
Some(daemon.socket()),
"the stand-in socket daemon must be discoverable, or this test proves nothing"
);
body();
drop(daemon);
unsafe {
std::env::remove_var(crate::search_rpc::TRUSTY_SEARCH_SOCKET_ENV);
match previous {
Some(p) => std::env::set_var(DATA_DIR_OVERRIDE_ENV, p),
None => std::env::remove_var(DATA_DIR_OVERRIDE_ENV),
}
}
drop(guard);
let _ = fs::remove_dir_all(&data_dir);
std::thread::sleep(std::time::Duration::from_millis(250));
let tcp_contacted = !matches!(
listener.accept(),
Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock
);
tcp_contacted || calls.load(Ordering::SeqCst) > 0
}
#[test]
fn ensure_project_indexed_never_writes_to_a_daemon_under_test() {
let root = scratch_dir("4255-ensure");
fs::create_dir_all(&root).expect("create fixture root");
let mut id = None;
let contacted = daemon_was_contacted_during(|| {
id = ensure_project_indexed(&root, true);
});
assert!(
!contacted,
"ensure_project_indexed contacted a live trusty-search daemon from a test \
process — that is the issue #4255 registry leak"
);
assert!(
id.is_none(),
"the guard suppressed the write, so no index was registered — handing \
back a pinnable id anyway is the #5091 fail-open shape"
);
let _ = fs::remove_dir_all(&root);
}
#[test]
fn index_files_inner_never_writes_to_a_daemon_under_test() {
let root = scratch_dir("4255-incremental");
fs::create_dir_all(&root).expect("create fixture root");
let file = root.join("fixture.rs");
fs::write(&file, "fn fixture() {}\n").expect("write fixture file");
let contacted = daemon_was_contacted_during(|| {
index_files_inner(&root, std::slice::from_ref(&file));
});
assert!(
!contacted,
"index_files_inner pushed fixture content to a live trusty-search daemon \
from a test process (issue #4255)"
);
let _ = fs::remove_dir_all(&root);
}
#[test]
fn index_options_builders_match_field_construction() {
assert_eq!(
IndexOptions::default().with_skip_vector(true),
IndexOptions {
allow_sensitive_path: false,
skip_vector: true,
}
);
assert_eq!(
IndexOptions::default().with_allow_sensitive_path(true),
IndexOptions {
allow_sensitive_path: true,
skip_vector: false,
}
);
assert_eq!(
IndexOptions::default()
.with_skip_vector(true)
.with_allow_sensitive_path(true),
IndexOptions {
allow_sensitive_path: true,
skip_vector: true,
}
);
}
use crate::search_rpc::{
CODE_CONFLICT, METHOD_INDEX_CREATE, METHOD_INDEX_REINDEX, METHOD_INDEX_STATUS,
METHOD_INDEXES_LIST,
};
use crate::uds_mock::{self, RpcError};
use std::sync::{Arc, Mutex};
type CallLog = Arc<Mutex<Vec<(String, serde_json::Value)>>>;
fn call_log() -> CallLog {
Arc::new(Mutex::new(Vec::new()))
}
fn calls(log: &CallLog) -> Vec<(String, serde_json::Value)> {
log.lock().unwrap_or_else(|e| e.into_inner()).clone()
}
fn methods(log: &CallLog) -> Vec<String> {
calls(log).into_iter().map(|(method, _)| method).collect()
}
fn params_of(log: &CallLog, method: &str, nth: usize) -> Option<serde_json::Value> {
calls(log)
.into_iter()
.filter(|(m, _)| m == method)
.map(|(_, p)| p)
.nth(nth)
}
fn recording(
log: CallLog,
answer: impl Fn(&str, usize) -> Result<serde_json::Value, RpcError> + Send + Sync + 'static,
) -> impl Fn(&str, serde_json::Value) -> uds_mock::MockFuture + Send + Sync + 'static {
move |method, params| {
let nth = {
let mut seen = log.lock().unwrap_or_else(|e| e.into_inner());
seen.push((method.to_string(), params));
seen.iter().filter(|(m, _)| m == method).count() - 1
};
let out = answer(method, nth);
Box::pin(async move { out })
}
}
fn reindex_lane(method: &str) -> Result<serde_json::Value, RpcError> {
if method == METHOD_INDEX_STATUS {
return Ok(serde_json::json!({ "chunk_count": 0 }));
}
if method == METHOD_INDEX_REINDEX {
return Ok(serde_json::json!({ "queued": true }));
}
Err(RpcError::new(
-32601,
format!("unexpected method: {method}"),
))
}
fn created_at(root: &Path) -> serde_json::Value {
serde_json::json!({
"id": "whatever",
"created": false,
"reason": "already exists",
"root_path": root.to_string_lossy(),
})
}
fn with_socket_daemon<T>(
tag: &str,
project_name: &str,
handler: impl Fn(&str, serde_json::Value) -> uds_mock::MockFuture + Send + Sync + 'static,
body: impl FnOnce(&Path) -> T,
) -> T {
let _guard = crate::data_dir::ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let data_dir = scratch_dir(&format!("7237-data-{tag}"));
fs::create_dir_all(&data_dir).unwrap();
let socket_dir = tempfile::tempdir().expect("tempdir for the mock socket");
let socket = socket_dir.path().join("s.sock");
unsafe {
std::env::set_var(crate::data_dir::DATA_DIR_OVERRIDE_ENV, &data_dir);
std::env::set_var(crate::test_harness::ALLOW_PRODUCTION_ENV, "1");
std::env::set_var(crate::search_rpc::TRUSTY_SEARCH_SOCKET_ENV, &socket);
}
let daemon = uds_mock::spawn_blocking_at(socket, handler);
let workspace = scratch_dir(&format!("7237-ws-{tag}"));
let project = workspace.join(project_name);
fs::create_dir_all(project.join(".git")).unwrap();
let out = body(&project);
drop(daemon);
unsafe {
std::env::remove_var(crate::search_rpc::TRUSTY_SEARCH_SOCKET_ENV);
std::env::remove_var(crate::data_dir::DATA_DIR_OVERRIDE_ENV);
std::env::remove_var(crate::test_harness::ALLOW_PRODUCTION_ENV);
}
let _ = fs::remove_dir_all(&workspace);
let _ = fs::remove_dir_all(&data_dir);
out
}
#[test]
fn the_create_call_carries_the_bare_create_index_params() {
let log = call_log();
let seen = Arc::clone(&log);
let (report, expected) = with_socket_daemon(
"bare-params",
"wire-project",
recording(seen, |method, _nth| {
if method == METHOD_INDEX_CREATE {
return Ok(serde_json::json!({ "id": "wire-project", "created": true }));
}
reindex_lane(method)
}),
|project| {
(
ensure_project_indexed_reporting(project, IndexOptions::default()),
crate::derive_index_id(project),
)
},
);
assert_eq!(
methods(&log).first().map(String::as_str),
Some(METHOD_INDEX_CREATE),
"the first thing the registration does must be the create: {:?}",
methods(&log)
);
let params = params_of(&log, METHOD_INDEX_CREATE, 0).expect("the create must have been sent");
assert_eq!(
params.get("id").and_then(serde_json::Value::as_str),
Some(expected.as_str()),
"the derived id rides on `id`, at the top level: {params}"
);
assert!(
params
.get("root_path")
.is_some_and(serde_json::Value::is_string),
"`root_path` must be a top-level string: {params}"
);
assert!(
params
.get("allow_sensitive_path")
.is_some_and(serde_json::Value::is_boolean),
"`allow_sensitive_path` must be a top-level bool: {params}"
);
assert!(
params
.get("skip_vector")
.is_some_and(serde_json::Value::is_boolean),
"`skip_vector` must be a top-level bool: {params}"
);
assert!(
params.get("index_id").is_none() && params.get("body").is_none(),
"the create is a REGISTRY-level write and must not use the index-scoped \
envelope, which the daemon would refuse as invalid_params: {params}"
);
assert_eq!(report.registration, IndexRegistration::Confirmed);
assert_eq!(report.index_id, Some(expected));
}
#[test]
fn a_confirmed_registration_triggers_a_reindex_over_the_socket() {
let log = call_log();
let seen = Arc::clone(&log);
with_socket_daemon(
"reindex",
"reindex-project",
recording(seen, |method, _nth| {
if method == METHOD_INDEX_CREATE {
return Ok(serde_json::json!({ "created": true }));
}
reindex_lane(method)
}),
|project| ensure_project_indexed_reporting(project, IndexOptions::default()),
);
assert_eq!(
methods(&log),
vec![
METHOD_INDEX_CREATE.to_string(),
METHOD_INDEX_STATUS.to_string(),
METHOD_INDEX_REINDEX.to_string(),
],
"create, then the freshness probe, then the trigger"
);
assert_eq!(
params_of(&log, METHOD_INDEX_REINDEX, 0).and_then(|p| {
p.get("index_id")
.and_then(serde_json::Value::as_str)
.map(str::to_string)
}),
Some("reindex-project".to_string()),
"the reindex names the index that was just registered"
);
}
#[test]
fn ensure_project_indexed_sends_allow_sensitive_path_through_to_create_body() {
for allow in [true, false] {
let log = call_log();
let seen = Arc::clone(&log);
with_socket_daemon(
&format!("wire-{allow}"),
"wire-project",
recording(seen, |method, _nth| {
if method == METHOD_INDEX_CREATE {
return Ok(serde_json::json!({ "created": true }));
}
reindex_lane(method)
}),
|project| ensure_project_indexed(project, allow),
);
let params =
params_of(&log, METHOD_INDEX_CREATE, 0).expect("the create must have been sent");
assert_eq!(
params.get("allow_sensitive_path"),
Some(&serde_json::Value::Bool(allow)),
"the create params must carry allow_sensitive_path={allow} all the way \
from ensure_project_indexed's parameter; got {params}"
);
}
}
#[test]
fn create_rejected_by_the_daemon_withholds_the_pinnable_id() {
fn refusing(method: &str, _nth: usize) -> Result<serde_json::Value, RpcError> {
if method == METHOD_INDEX_CREATE {
return Err(RpcError::internal("the corpus would not open"));
}
reindex_lane(method)
}
let id = with_socket_daemon(
"refuse-a",
"refused",
recording(call_log(), refusing),
|p| ensure_project_indexed(p, false),
);
assert_eq!(
id, None,
"ensure_project_indexed returned a pinnable id after the daemon REFUSED \
the create — pinning it makes every later search 404 (#5091)"
);
let id = with_socket_daemon(
"refuse-b",
"refused",
recording(call_log(), refusing),
|p| ensure_project_indexed_with(p, IndexOptions::default().with_skip_vector(true)),
);
assert_eq!(
id, None,
"ensure_project_indexed_with returned a pinnable id after the daemon \
REFUSED the create — #5091"
);
let (report, expected) = with_socket_daemon(
"refuse-c",
"refused",
recording(call_log(), refusing),
|p| {
(
ensure_project_indexed_reporting(p, IndexOptions::default()),
crate::derive_index_id(p),
)
},
);
assert_eq!(
report.registration,
IndexRegistration::NotConfirmed,
"a refused create is not a registration"
);
assert_eq!(
report.index_id,
Some(expected),
"the derived id stays available for logging and GC — it is the PIN that \
is withheld, not the id"
);
}
#[test]
fn a_daemon_refusal_is_not_a_registration() {
let log = call_log();
let seen = Arc::clone(&log);
let outcome = with_socket_daemon(
"refusal-arm",
"refusal-arm",
recording(seen, |_method, _nth| {
Err(RpcError::internal("the embedder never initialised"))
}),
|project| {
let socket = crate::search_rpc::search_socket().expect("derive the socket");
best_effort_create_index(&socket, "api", project, IndexOptions::default())
},
);
assert_eq!(outcome, CreateOutcome::NotConfirmed);
assert_eq!(
methods(&log),
vec![METHOD_INDEX_CREATE.to_string()],
"an unrecoverable refusal must not go on to read the registry"
);
}
#[test]
fn a_malformed_create_reply_is_not_a_registration() {
use std::io::{BufRead, BufReader, Write};
use std::os::unix::fs::PermissionsExt;
let dir = tempfile::tempdir().expect("tempdir for the garbage daemon");
let socket = dir.path().join("s.sock");
let listener =
std::os::unix::net::UnixListener::bind(&socket).expect("bind the garbage daemon");
std::fs::set_permissions(&socket, std::fs::Permissions::from_mode(0o600))
.expect("harden the garbage socket");
std::fs::set_permissions(dir.path(), std::fs::Permissions::from_mode(0o700))
.expect("harden the garbage socket directory");
let server = std::thread::spawn(move || {
let Ok((stream, _)) = listener.accept() else {
return;
};
let mut reader = BufReader::new(stream);
let mut request = String::new();
let _ = reader.read_line(&mut request);
let mut stream = reader.into_inner();
let _ = stream.write_all(b"this is not a json-rpc frame\n");
let _ = stream.flush();
});
let root = scratch_dir("7237-garbage");
fs::create_dir_all(&root).unwrap();
let outcome = best_effort_create_index(&socket, "api", &root, IndexOptions::default());
let _ = server.join();
let _ = fs::remove_dir_all(&root);
assert_eq!(
outcome,
CreateOutcome::NotConfirmed,
"a reply that cannot be read acknowledges nothing"
);
}
#[test]
fn a_missing_socket_registers_nothing_and_contacts_no_tcp_port() {
let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("bind the legacy port");
let addr = listener.local_addr().expect("legacy port local_addr");
listener
.set_nonblocking(true)
.expect("legacy port set_nonblocking");
let _guard = crate::data_dir::ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let data_dir = scratch_dir("7237-nosocket-data");
fs::create_dir_all(&data_dir).unwrap();
unsafe {
std::env::set_var(crate::data_dir::DATA_DIR_OVERRIDE_ENV, &data_dir);
std::env::set_var(crate::test_harness::ALLOW_PRODUCTION_ENV, "1");
}
crate::write_daemon_addr("trusty-search", &addr.to_string()).expect("publish the legacy addr");
assert_eq!(
crate::resolve_daemon_base_url("trusty-search"),
Some(format!("http://{addr}")),
"the legacy discovery file must resolve, or this test proves nothing"
);
let socket = crate::search_rpc::search_socket().expect("derive the socket path");
assert!(
!socket.exists(),
"the rig's premise is that nothing is listening on the socket"
);
let project = scratch_dir("7237-nosocket");
fs::create_dir_all(project.join(".git")).unwrap();
let report = ensure_project_indexed_reporting(&project, IndexOptions::default());
unsafe {
std::env::remove_var(crate::data_dir::DATA_DIR_OVERRIDE_ENV);
std::env::remove_var(crate::test_harness::ALLOW_PRODUCTION_ENV);
}
let _ = fs::remove_dir_all(&project);
let _ = fs::remove_dir_all(&data_dir);
assert_eq!(
report.registration,
IndexRegistration::DaemonUnreachable,
"no socket means nothing was sent — that is not a registration"
);
assert_eq!(
pinnable_index_id(report),
None,
"an unconfirmed registration must not advance a caller's pin (#5091)"
);
std::thread::sleep(std::time::Duration::from_millis(250));
assert!(
matches!(listener.accept(), Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock),
"the registration dialled the retired loopback port — the #7237 fallback \
that must not exist"
);
}
#[test]
fn registered_root_from_response_reads_the_already_exists_root() {
let body = serde_json::json!({
"id": "api", "created": false, "reason": "already exists", "root_path": "/srv/other"
});
assert_eq!(
registered_root_from_response(&body),
Some("/srv/other".to_string())
);
}
#[test]
fn registered_root_from_response_ignores_a_fresh_create() {
let body = serde_json::json!({ "id": "api", "created": true, "root_path": "/srv/api" });
assert_eq!(registered_root_from_response(&body), None);
}
#[test]
fn registered_root_from_response_tolerates_a_daemon_that_omits_it() {
assert_eq!(
registered_root_from_response(&serde_json::json!({ "id": "api", "created": false })),
None
);
assert_eq!(
registered_root_from_response(&serde_json::Value::String("not an object".into())),
None
);
assert_eq!(
registered_root_from_response(&serde_json::Value::Null),
None
);
assert_eq!(
registered_root_from_response(&serde_json::json!({ "created": false, "root_path": 42 })),
None,
"a non-string root_path must yield None, not a panic"
);
}
#[test]
fn create_index_response_for_a_different_tree_reports_a_conflict() {
let registered = scratch_dir("mismatch-registered");
fs::create_dir_all(®istered).unwrap();
let answer = created_at(®istered);
let outcome = with_socket_daemon(
"mismatch",
"mismatch-requested",
recording(call_log(), move |method, _nth| {
if method == METHOD_INDEX_CREATE {
return Ok(answer.clone());
}
reindex_lane(method)
}),
|project| {
let socket = crate::search_rpc::search_socket().expect("derive the socket");
best_effort_create_index(&socket, "api", project, IndexOptions::default())
},
);
assert_eq!(
outcome,
CreateOutcome::Conflict { existing_id: None },
"an answer naming a different tree must not confirm the registration"
);
let _ = fs::remove_dir_all(®istered);
}
#[test]
fn create_index_response_for_the_same_tree_is_confirmed() {
let outcome = with_socket_daemon(
"same-tree",
"same-tree",
|method, params| {
let answer = if method == METHOD_INDEX_CREATE {
let root = params
.get("root_path")
.cloned()
.unwrap_or(serde_json::Value::Null);
Ok(serde_json::json!({
"id": "api", "created": false, "reason": "already exists", "root_path": root
}))
} else {
reindex_lane(method)
};
Box::pin(async move { answer })
},
|project| {
let socket = crate::search_rpc::search_socket().expect("derive the socket");
best_effort_create_index(&socket, "api", project, IndexOptions::default())
},
);
assert_eq!(outcome, CreateOutcome::Confirmed);
}
fn root_mismatch_message(index_id: &str, registered: &Path) -> String {
format!(
"index '{index_id}' is registered at {:?}; it cannot be re-registered \
because one index identifies one directory tree",
registered.display()
)
}
fn root_collision_message(root: &Path, existing_id: &str) -> String {
format!(
"root_path {:?} is already registered to index '{existing_id}'; two \
indexes cannot share one on-disk corpus (issues #2305, #2336)",
root.display()
)
}
#[test]
fn registration_matches_an_existing_index_by_root_path() {
let other = scratch_dir("6864-other-checkout");
fs::create_dir_all(&other).unwrap();
let elsewhere = other.clone();
let mine: Arc<Mutex<Option<PathBuf>>> = Arc::new(Mutex::new(None));
let recorded = Arc::clone(&mine);
let report = with_socket_daemon(
"match",
"trusty-tools",
move |method, params| {
let answer = if method == METHOD_INDEX_CREATE {
let requested = params
.get("root_path")
.and_then(serde_json::Value::as_str)
.unwrap_or_default()
.to_string();
*recorded.lock().unwrap_or_else(|e| e.into_inner()) =
Some(PathBuf::from(requested));
Err(RpcError::new(
CODE_CONFLICT,
root_mismatch_message("trusty-tools", &elsewhere),
))
} else if method == METHOD_INDEXES_LIST {
let requested = recorded
.lock()
.unwrap_or_else(|e| e.into_inner())
.clone()
.unwrap_or_default();
Ok(serde_json::json!({
"indexes": [
{ "id": "trusty-tools", "root_path": elsewhere.to_string_lossy() },
{
"id": "trusty-tools-checkout",
"root_path": requested.to_string_lossy(),
},
]
}))
} else {
reindex_lane(method)
};
Box::pin(async move { answer })
},
|project| ensure_project_indexed_reporting(project, IndexOptions::default()),
);
assert_eq!(
report.registration,
IndexRegistration::Confirmed,
"an index registered at this root IS a confirmed registration (#6864)"
);
assert_eq!(
report.index_id,
Some("trusty-tools-checkout".to_string()),
"the report must carry the id that serves this tree, not the colliding \
basename the daemon refused (#6864)"
);
let _ = fs::remove_dir_all(&other);
}
#[test]
fn a_root_path_collision_recovers_the_owning_index_without_a_registry_read() {
let log = call_log();
let seen = Arc::clone(&log);
let report = with_socket_daemon(
"collision",
"trusty-tools",
recording(seen, |method, _nth| {
if method == METHOD_INDEX_CREATE {
return Err(RpcError::new(
CODE_CONFLICT,
root_collision_message(
Path::new("/nonexistent/checkout/trusty-tools"),
"trusty-tools-checkout",
),
));
}
reindex_lane(method)
}),
|project| ensure_project_indexed_reporting(project, IndexOptions::default()),
);
assert_eq!(report.registration, IndexRegistration::Confirmed);
assert_eq!(
report.index_id,
Some("trusty-tools-checkout".to_string()),
"the id the refusal named is the one to pin (#6864)"
);
assert!(
!methods(&log).contains(&METHOD_INDEXES_LIST.to_string()),
"the refusal already named the owning index, so the registry read is \
wasted work: {:?}",
methods(&log)
);
}
#[test]
fn registration_falls_back_to_a_collision_resistant_id() {
let other = scratch_dir("6864-fallback-other");
fs::create_dir_all(&other).unwrap();
let elsewhere = other.clone();
let log = call_log();
let seen = Arc::clone(&log);
let (report, expected) = with_socket_daemon(
"fallback",
"trusty-tools",
recording(seen, move |method, nth| {
if method == METHOD_INDEX_CREATE {
if nth == 0 {
return Err(RpcError::new(
CODE_CONFLICT,
root_mismatch_message("trusty-tools", &elsewhere),
));
}
return Ok(serde_json::json!({ "created": true }));
}
if method == METHOD_INDEXES_LIST {
return Ok(serde_json::json!({
"indexes": [
{ "id": "trusty-tools", "root_path": elsewhere.to_string_lossy() }
]
}));
}
reindex_lane(method)
}),
|project| {
(
ensure_project_indexed_reporting(project, IndexOptions::default()),
crate::derive_checkout_index_id(project),
)
},
);
assert_eq!(
report.registration,
IndexRegistration::Confirmed,
"the fallback create landed, so the registration is confirmed (#6864)"
);
assert_eq!(
report.index_id, expected,
"the id must be the shared checkout derivation, not a scheme invented here \
(#6149 / #6864)"
);
assert_eq!(
calls(&log)
.iter()
.filter(|(m, _)| m == METHOD_INDEX_CREATE)
.count(),
2,
"exactly one retry, under the digest id"
);
let _ = fs::remove_dir_all(&other);
}
const LATE_REGISTERED_ID: &str = "late-create-project-checkout";
#[test]
fn a_late_create_is_confirmed_by_polling_the_registry() {
let log = call_log();
let recorder = Arc::clone(&log);
let registered: Arc<Mutex<Option<String>>> = Arc::new(Mutex::new(None));
let daemon_side = Arc::clone(®istered);
let report = with_socket_daemon(
"late-create",
"late-create-project",
move |method, params| {
recorder
.lock()
.unwrap_or_else(|e| e.into_inner())
.push((method.to_string(), params.clone()));
let registered = Arc::clone(&daemon_side);
let method = method.to_string();
Box::pin(async move {
if method == METHOD_INDEX_CREATE {
tokio::time::sleep(Duration::from_millis(1_500)).await;
let root = params
.get("root_path")
.and_then(serde_json::Value::as_str)
.expect("the create carries a root_path")
.to_string();
*registered.lock().unwrap_or_else(|e| e.into_inner()) = Some(root);
return Ok(serde_json::json!({ "created": true }));
}
if method == METHOD_INDEXES_LIST {
let listed = registered.lock().unwrap_or_else(|e| e.into_inner()).clone();
return Ok(match listed {
Some(root) => serde_json::json!({
"indexes": [{ "id": LATE_REGISTERED_ID, "root_path": root }]
}),
None => serde_json::json!({ "indexes": [] }),
});
}
reindex_lane(&method)
})
},
|project| ensure_project_indexed_reporting(project, IndexOptions::default()),
);
assert_eq!(
report.registration,
IndexRegistration::Confirmed,
"the daemon registered the index; a client budget that elapsed first is not a \
refusal (#7237)"
);
assert_eq!(
report.index_id,
Some(LATE_REGISTERED_ID.to_string()),
"the pinned id must be the one the registry reports for this tree"
);
let seen = methods(&log);
assert_eq!(
seen.first().map(String::as_str),
Some(METHOD_INDEX_CREATE),
"the create still goes first: {seen:?}"
);
assert!(
seen.iter().any(|m| m == METHOD_INDEXES_LIST),
"the unanswered create must be settled by reading the registry: {seen:?}"
);
assert_eq!(
seen.last().map(String::as_str),
Some(METHOD_INDEX_REINDEX),
"a confirmed registration still gets its reindex trigger: {seen:?}"
);
}
#[test]
fn an_unconfirmed_registration_fires_no_reindex() {
let log = call_log();
let seen = Arc::clone(&log);
let report = with_socket_daemon(
"no-reindex",
"no-reindex-project",
recording(seen, |method, _nth| {
if method == METHOD_INDEX_CREATE {
return Err(RpcError::internal("the corpus would not open"));
}
reindex_lane(method)
}),
|project| ensure_project_indexed_reporting(project, IndexOptions::default()),
);
assert_eq!(report.registration, IndexRegistration::NotConfirmed);
assert_eq!(
pinnable_index_id(report),
None,
"an index the daemon never registered must stay unpinned (#5091)"
);
assert_eq!(
methods(&log),
vec![METHOD_INDEX_CREATE.to_string()],
"nothing may be reindexed against an index the daemon has not registered"
);
}
#[test]
fn a_create_the_daemon_hung_up_on_is_confirmed_by_polling_the_registry() {
let log = call_log();
let recorder = Arc::clone(&log);
let registered: Arc<Mutex<Option<String>>> = Arc::new(Mutex::new(None));
let daemon_side = Arc::clone(®istered);
let report = with_socket_daemon(
"hung-up-create",
"hung-up-create-project",
move |method, params| {
recorder
.lock()
.unwrap_or_else(|e| e.into_inner())
.push((method.to_string(), params.clone()));
let registered = Arc::clone(&daemon_side);
let method = method.to_string();
Box::pin(async move {
if method == METHOD_INDEX_CREATE {
let root = params
.get("root_path")
.and_then(serde_json::Value::as_str)
.expect("the create carries a root_path")
.to_string();
*registered.lock().unwrap_or_else(|e| e.into_inner()) = Some(root);
panic!("the mock daemon dropped the create connection after registering");
}
if method == METHOD_INDEXES_LIST {
let listed = registered.lock().unwrap_or_else(|e| e.into_inner()).clone();
return Ok(match listed {
Some(root) => serde_json::json!({
"indexes": [{ "id": LATE_REGISTERED_ID, "root_path": root }]
}),
None => serde_json::json!({ "indexes": [] }),
});
}
reindex_lane(&method)
})
},
|project| ensure_project_indexed_reporting(project, IndexOptions::default()),
);
assert_eq!(
report.registration,
IndexRegistration::Confirmed,
"a daemon that hung up after registering has still registered (#7237)"
);
assert_eq!(
report.index_id,
Some(LATE_REGISTERED_ID.to_string()),
"the pinned id must be the one the registry reports for this tree"
);
let seen = methods(&log);
assert!(
seen.iter().any(|m| m == METHOD_INDEXES_LIST),
"a hang-up must be settled by reading the registry: {seen:?}"
);
}