use std::path::Path;
use std::time::Duration;
use super::{
CREATE_TIMEOUT, IndexOptions, IndexRegistration, best_effort_create_index,
registered_root_from_response,
};
use crate::search_rpc::{self, SearchRpcError};
use crate::uds::UdsRpcError;
const LIST_TIMEOUT: Duration = CREATE_TIMEOUT;
const COLLISION_MARKER: &str = "is already registered to index '";
const OPERATOR_ADVICE_MARKER: &str = ". Use the relocate endpoint";
fn without_operator_advice(message: &str) -> &str {
message
.split_once(OPERATOR_ADVICE_MARKER)
.map_or(message, |(facts, _)| facts)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(super) enum CreateOutcome {
Confirmed,
Conflict { existing_id: Option<String> },
NotConfirmed,
Unanswered,
}
pub(super) fn classify_create_result(
result: &serde_json::Value,
index_id: &str,
root: &Path,
root_display: &str,
) -> CreateOutcome {
match registered_root_from_response(result) {
Some(registered) if !crate::identifies_same_path(Path::new(®istered), root) => {
tracing::warn!(
"trusty-search index '{index_id}' is registered at {registered}, not at the \
requested {root_display}; withholding confirmation so the caller cannot pin \
an index that searches a different tree"
);
CreateOutcome::Conflict { existing_id: None }
}
_ => {
tracing::debug!("registered trusty-search index '{index_id}' (root={root_display})");
CreateOutcome::Confirmed
}
}
}
pub(super) fn classify_create_failure(
err: &anyhow::Error,
index_id: &str,
root_display: &str,
) -> CreateOutcome {
let Some(refusal) = err.downcast_ref::<SearchRpcError>() else {
if create_left_unanswered(err) {
tracing::info!(
"trusty-search has not answered the index registration for '{index_id}' at \
{root_display} within the create budget ({err:#}); checking the registry for \
a late registration (#7237)"
);
return CreateOutcome::Unanswered;
}
tracing::warn!("trusty-search index registration for '{index_id}' failed: {err:#}");
return CreateOutcome::NotConfirmed;
};
if refusal.is_conflict() {
tracing::info!(
"trusty-search index registration for '{index_id}' at {root_display} \
conflicts ({}); resolving it to the index that already serves that tree",
without_operator_advice(&refusal.message)
);
return CreateOutcome::Conflict {
existing_id: existing_id_from_conflict(&refusal.message),
};
}
tracing::warn!("trusty-search index registration for '{index_id}' was refused: {refusal}");
CreateOutcome::NotConfirmed
}
fn create_left_unanswered(err: &anyhow::Error) -> bool {
matches!(
err.downcast_ref::<UdsRpcError>(),
Some(
UdsRpcError::Timeout { .. }
| UdsRpcError::NoResponse { .. }
| UdsRpcError::Read { .. }
| UdsRpcError::HalfClose { .. }
)
)
}
pub(super) fn create_and_reconcile(
socket: &Path,
index_id: &str,
root: &Path,
opts: IndexOptions,
) -> (String, IndexRegistration) {
match best_effort_create_index(socket, index_id, root, opts) {
CreateOutcome::Confirmed => (index_id.to_string(), IndexRegistration::Confirmed),
CreateOutcome::NotConfirmed => (index_id.to_string(), IndexRegistration::NotConfirmed),
CreateOutcome::Unanswered => {
match super::confirm::confirm_after_no_answer(socket, index_id, root) {
Some(resolved) => (resolved, IndexRegistration::Confirmed),
None => (index_id.to_string(), IndexRegistration::NotConfirmed),
}
}
CreateOutcome::Conflict { existing_id } => {
match resolve_colliding_id(socket, index_id, root, opts, existing_id) {
Some(resolved) => (resolved, IndexRegistration::Confirmed),
None => (index_id.to_string(), IndexRegistration::NotConfirmed),
}
}
}
}
fn resolve_colliding_id(
socket: &Path,
derived: &str,
root: &Path,
opts: IndexOptions,
existing_id: Option<String>,
) -> Option<String> {
if let Some(id) = existing_id {
tracing::info!(
"trusty-search already serves {} as index '{id}'; pinning that instead of \
the derived '{derived}' (#6864)",
root.display()
);
return Some(id);
}
if let Some(body) = fetch_index_list(socket, ListFailure::Warn)
&& let Some(id) = index_id_serving_root(&body, root)
{
tracing::info!(
"trusty-search index '{derived}' identifies another tree; {} is registered \
as '{id}' and that is what this session pins (#6864)",
root.display()
);
return Some(id);
}
let fresh = crate::derive_checkout_index_id(root)?;
tracing::info!(
"no trusty-search index is registered for {}; registering it under the \
collision-resistant id '{fresh}' because '{derived}' names another tree (#6864)",
root.display()
);
let resolved = match best_effort_create_index(socket, &fresh, root, opts) {
CreateOutcome::Confirmed => Some(fresh),
CreateOutcome::Conflict { existing_id } => existing_id,
CreateOutcome::NotConfirmed => None,
CreateOutcome::Unanswered => super::confirm::confirm_after_no_answer(socket, &fresh, root),
};
if let Some(id) = &resolved {
tracing::info!(
"trusty-search serves {} as index '{id}' (#6864)",
root.display()
);
}
resolved
}
pub(super) fn existing_id_from_conflict(message: &str) -> Option<String> {
let (_, rest) = message.split_once(COLLISION_MARKER)?;
let (id, _) = rest.split_once('\'')?;
(!id.is_empty()).then(|| id.to_string())
}
pub(super) fn index_id_serving_root(body: &serde_json::Value, root: &Path) -> Option<String> {
for entry in body.get("indexes")?.as_array()? {
let Some(registered) = entry.get("root_path").and_then(|v| v.as_str()) else {
continue;
};
if !crate::identifies_same_path(Path::new(registered), root) {
continue;
}
if let Some(id) = entry.get("id").and_then(|v| v.as_str()) {
return Some(id.to_string());
}
}
None
}
pub(super) fn fetch_index_list(
socket: &Path,
on_failure: ListFailure,
) -> Option<serde_json::Value> {
match search_rpc::call_blocking(
socket,
search_rpc::METHOD_INDEXES_LIST,
serde_json::json!({ "details": true }),
LIST_TIMEOUT,
) {
Ok(body) => Some(body),
Err(e) => {
match on_failure {
ListFailure::Warn => tracing::warn!("trusty-search index list failed: {e:#}"),
ListFailure::Quiet => tracing::debug!("trusty-search index list failed: {e:#}"),
}
None
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) enum ListFailure {
Warn,
Quiet,
}
#[cfg(test)]
mod tests {
use super::*;
use crate::log_buffer::capture_logs;
use crate::uds_mock;
use std::sync::{Arc, Mutex};
fn collision_message(root: &str, existing_id: &str) -> String {
format!(
"root_path {root:?} is already registered to index '{existing_id}'; two \
indexes cannot share one on-disk corpus (issues #2305, #2336)"
)
}
fn mismatch_message(index_id: &str, registered: &str, requested: &str) -> String {
format!(
"index '{index_id}' is registered at {registered:?}; it cannot be re-registered \
at {requested:?} because one index identifies one directory tree. Use the \
relocate endpoint to move it, or register the other tree under a distinct id"
)
}
fn conflict_error(message: String) -> anyhow::Error {
anyhow::Error::new(SearchRpcError {
method: search_rpc::METHOD_INDEX_CREATE.to_string(),
code: search_rpc::CODE_CONFLICT,
message,
})
}
#[test]
fn a_resolvable_conflict_is_logged_at_info_without_operator_advice() {
let registered = "/Users/masa/Projects/trusty-tools";
let requested = "/Users/masa/trusty-mpm-projects/bobmatnyc/trusty-tools";
let err = conflict_error(mismatch_message("trusty-tools", registered, requested));
let (outcome, lines) =
capture_logs(|| classify_create_failure(&err, "trusty-tools", requested));
assert_eq!(outcome, CreateOutcome::Conflict { existing_id: None });
assert_eq!(lines.len(), 1, "expected one line, got {lines:?}");
let line = &lines[0];
assert!(
line.contains("INFO") && !line.contains("WARN"),
"a conflict this code resolves itself is not a warning: {line}"
);
assert!(line.contains(registered), "line was: {line}");
assert!(line.contains(requested), "line was: {line}");
assert!(
!line.contains("Use the relocate endpoint"),
"the daemon's operator advice does not belong in a self-healing line: {line}"
);
}
#[test]
fn an_unrecoverable_refusal_still_warns() {
let err = anyhow::Error::new(SearchRpcError {
method: search_rpc::METHOD_INDEX_CREATE.to_string(),
code: -32603,
message: "internal error".to_string(),
});
let (outcome, lines) =
capture_logs(|| classify_create_failure(&err, "trusty-tools", "/nonexistent/tree"));
assert_eq!(outcome, CreateOutcome::NotConfirmed);
assert_eq!(lines.len(), 1, "expected one line, got {lines:?}");
assert!(
lines[0].contains("WARN"),
"nothing recovers this, so it stays a warning: {}",
lines[0]
);
}
#[test]
fn a_collision_message_passes_through_whole() {
let message = collision_message("/Users/masa/checkout/trusty-tools", "trusty-tools");
assert_eq!(without_operator_advice(&message), message);
}
#[test]
fn an_unanswered_create_the_registry_confirms_emits_no_warning() {
let root = tempfile::tempdir().expect("tempdir for the indexed tree");
let ((resolved, registration), lines) = with_daemon(late_registering_daemon(), |socket| {
capture_logs(|| {
create_and_reconcile(socket, "asked-for", root.path(), IndexOptions::default())
})
});
assert_eq!(
registration,
IndexRegistration::Confirmed,
"the daemon registered the index; the confirm loop found it: {lines:?}"
);
assert_eq!(
resolved, LATE_REGISTERED_ID,
"the pinned id must be the one the registry names for this tree"
);
let warnings: Vec<&String> = lines.iter().filter(|l| l.contains("WARN")).collect();
assert!(
warnings.is_empty(),
"a registration that succeeded must warn about nothing, saw {warnings:?} in {lines:?}"
);
assert!(
lines.iter().any(|l| l.contains("INFO")),
"the recovery still reports itself, at info: {lines:?}"
);
}
const LATE_REGISTERED_ID: &str = "registered-late-by-the-daemon";
fn late_registering_daemon()
-> impl Fn(&str, serde_json::Value) -> uds_mock::MockFuture + Send + Sync + 'static {
let registered: Arc<Mutex<Option<String>>> = Arc::new(Mutex::new(None));
move |method, params| {
let registered = Arc::clone(®istered);
let method = method.to_string();
Box::pin(async move {
if method == search_rpc::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");
}
let listed = registered.lock().unwrap_or_else(|e| e.into_inner()).clone();
Ok(match listed {
Some(root) => serde_json::json!({
"indexes": [{ "id": LATE_REGISTERED_ID, "root_path": root }]
}),
None => serde_json::json!({ "indexes": [] }),
})
})
}
}
fn with_daemon<T>(
handler: impl Fn(&str, serde_json::Value) -> uds_mock::MockFuture + Send + Sync + 'static,
body: impl FnOnce(&Path) -> T,
) -> T {
let dir = tempfile::tempdir().expect("tempdir for the mock socket");
let socket = dir.path().join("s.sock");
let daemon = uds_mock::spawn_blocking_at(socket.clone(), handler);
let out = body(&socket);
drop(daemon);
out
}
#[test]
fn an_unanswered_create_is_not_a_refusal() {
let path = std::path::PathBuf::from("/nonexistent/trusty-search.sock");
let unanswered = [
UdsRpcError::NoResponse { path: path.clone() },
UdsRpcError::Timeout {
path: path.clone(),
timeout: Duration::from_secs(1),
},
UdsRpcError::Read {
path: path.clone(),
source: std::io::Error::other("truncated"),
},
];
for err in unanswered {
let rendered = err.to_string();
let wrapped = anyhow::Error::new(err).context("call search.index.create");
assert_eq!(
classify_create_failure(&wrapped, "writing", "/nonexistent/writing"),
CreateOutcome::Unanswered,
"the daemon said nothing, so this is not a refusal: {rendered}"
);
}
}
#[test]
fn a_dial_failure_is_not_an_unanswered_create() {
let path = std::path::PathBuf::from("/nonexistent/trusty-search.sock");
let decode = UdsRpcError::Decode {
path: path.clone(),
source: serde_json::from_str::<serde_json::Value>("not json").unwrap_err(),
};
assert_eq!(
classify_create_failure(
&anyhow::Error::new(decode),
"writing",
"/nonexistent/writing"
),
CreateOutcome::NotConfirmed,
"a reply that cannot be decoded is an answer, not a silence"
);
let plain = anyhow::anyhow!("the trusty-search search.index.create worker thread panicked");
assert_eq!(
classify_create_failure(&plain, "writing", "/nonexistent/writing"),
CreateOutcome::NotConfirmed,
"an error carrying no transport variant confirms nothing and polls nothing"
);
}
#[test]
fn a_root_collision_message_carries_the_existing_id() {
assert_eq!(
existing_id_from_conflict(&collision_message(
"/Users/masa/checkout/trusty-tools",
"trusty-tools-checkout"
)),
Some("trusty-tools-checkout".to_string())
);
}
#[test]
fn a_root_mismatch_message_names_no_existing_index() {
let mismatch = "index 'trusty-tools' is registered at \
\"/Users/masa/Projects/trusty-tools\"; it cannot be re-registered at \
\"/Users/masa/checkout/trusty-tools\" because one index identifies one \
directory tree";
assert_eq!(existing_id_from_conflict(mismatch), None);
assert_eq!(existing_id_from_conflict(""), None);
assert_eq!(
existing_id_from_conflict("root_path is already registered to index '"),
None,
"an unterminated quote names no id"
);
assert_eq!(
existing_id_from_conflict("root_path is already registered to index ''"),
None,
"an empty id is not an id"
);
}
#[test]
fn index_id_serving_root_matches_on_the_tree_not_the_id() {
let mine = std::env::temp_dir();
let body = serde_json::json!({
"indexes": [
{ "id": "trusty-tools", "root_path": "/nonexistent/other/trusty-tools" },
{ "id": "trusty-tools-checkout", "root_path": mine.to_string_lossy() },
]
});
assert_eq!(
index_id_serving_root(&body, &mine),
Some("trusty-tools-checkout".to_string()),
"the entry rooted at this tree is the one to pin, whatever its id"
);
}
#[test]
fn index_id_serving_root_is_none_when_no_entry_matches() {
let body = serde_json::json!({
"indexes": [{ "id": "api", "root_path": "/nonexistent/work/api" }]
});
assert_eq!(
index_id_serving_root(&body, Path::new("/nonexistent/work/other")),
None
);
}
#[test]
fn index_id_serving_root_tolerates_a_malformed_body() {
let root = std::env::temp_dir();
assert_eq!(
index_id_serving_root(&serde_json::Value::String("nope".into()), &root),
None
);
assert_eq!(index_id_serving_root(&serde_json::json!({}), &root), None);
assert_eq!(
index_id_serving_root(&serde_json::json!({ "indexes": "nope" }), &root),
None
);
assert_eq!(
index_id_serving_root(
&serde_json::json!({ "indexes": [{ "id": "x", "root_path": null }] }),
&root
),
None
);
assert_eq!(
index_id_serving_root(
&serde_json::json!({ "indexes": [{ "root_path": root.to_string_lossy() }] }),
&root
),
None,
"an entry with no id cannot be pinned"
);
}
}