use std::path::Path;
use std::time::{Duration, Instant};
use super::reconcile::{ListFailure, fetch_index_list, index_id_serving_root};
const CONFIRM_DEADLINE: Duration = Duration::from_secs(4);
const CONFIRM_POLL_GAP: Duration = Duration::from_millis(250);
pub(super) fn confirm_after_no_answer(
socket: &Path,
index_id: &str,
root: &Path,
) -> Option<String> {
confirm_within(socket, index_id, root, CONFIRM_DEADLINE, CONFIRM_POLL_GAP)
}
fn confirm_within(
socket: &Path,
index_id: &str,
root: &Path,
deadline: Duration,
gap: Duration,
) -> Option<String> {
let started = Instant::now();
let mut reads: u32 = 0;
loop {
reads += 1;
if let Some(body) = fetch_index_list(socket, ListFailure::Quiet)
&& let Some(registered) = index_id_serving_root(&body, root)
{
tracing::info!(
"trusty-search registered {} as index '{registered}' after the create call \
went unanswered ({reads} registry read(s), {:?}); pinning it (#7237)",
root.display(),
started.elapsed()
);
return Some(registered);
}
let elapsed = started.elapsed();
if elapsed >= deadline {
break;
}
std::thread::sleep(gap.min(deadline - elapsed));
}
tracing::warn!(
"trusty-search has not registered {} after {reads} registry read(s) over {deadline:?}; \
it may still be registering '{index_id}' in the background, but nothing here retries \
and the id is withheld so nothing pins an index that may not exist (#7237)",
root.display()
);
None
}
#[cfg(test)]
mod tests {
use super::*;
use crate::uds_mock::{self, MockFuture, RpcError};
use std::sync::Arc;
use std::sync::atomic::{AtomicU32, Ordering};
fn with_daemon<T>(
handler: impl Fn(&str, serde_json::Value) -> 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
}
fn listing(id: &str, root: &Path) -> serde_json::Value {
serde_json::json!({
"indexes": [{ "id": id, "root_path": root.to_string_lossy() }]
})
}
#[test]
fn confirm_within_stops_at_the_deadline_and_withholds() {
let reads = Arc::new(AtomicU32::new(0));
let counter = Arc::clone(&reads);
let started = Instant::now();
let confirmed = with_daemon(
move |_method, _params| {
counter.fetch_add(1, Ordering::SeqCst);
Box::pin(async { Ok(serde_json::json!({ "indexes": [] })) })
},
|socket| {
confirm_within(
socket,
"never-registered",
Path::new("/nonexistent/never/registered"),
Duration::from_millis(120),
Duration::from_millis(20),
)
},
);
assert_eq!(
confirmed, None,
"an index the daemon never registers must stay unpinned (#5091)"
);
assert!(
reads.load(Ordering::SeqCst) > 1,
"the confirm must poll rather than read once, saw {} read(s)",
reads.load(Ordering::SeqCst)
);
assert!(
started.elapsed() < Duration::from_secs(2),
"the poll must stop at its deadline, took {:?}",
started.elapsed()
);
}
#[test]
fn an_exhausted_confirm_deadline_is_still_a_warning() {
let root = Path::new("/nonexistent/never/registered");
let (confirmed, lines) = with_daemon(
move |_method, _params| Box::pin(async { Ok(serde_json::json!({ "indexes": [] })) }),
|socket| {
crate::log_buffer::capture_logs(|| {
confirm_within(
socket,
"never-registered",
root,
Duration::from_millis(60),
Duration::from_millis(20),
)
})
},
);
assert_eq!(confirmed, None, "an empty registry confirms nothing");
assert_eq!(lines.len(), 1, "expected one line, got {lines:?}");
assert!(
lines[0].contains("WARN") && lines[0].contains(&root.display().to_string()),
"the give-up is the outcome an operator has to act on: {}",
lines[0]
);
}
#[test]
fn confirm_within_never_confirms_an_index_at_another_root() {
let confirmed = with_daemon(
move |_method, _params| {
let body = listing("someone-else", Path::new("/nonexistent/other/tree"));
Box::pin(async move { Ok(body) })
},
|socket| {
confirm_within(
socket,
"mine",
Path::new("/nonexistent/my/tree"),
Duration::from_millis(60),
Duration::from_millis(20),
)
},
);
assert_eq!(
confirmed, None,
"an index registered at a DIFFERENT tree is not this registration"
);
}
#[test]
fn confirm_within_treats_a_refusing_registry_as_no_answer() {
let confirmed = with_daemon(
move |_method, _params| {
Box::pin(async { Err(RpcError::internal("the registry is unavailable")) })
},
|socket| {
confirm_within(
socket,
"mine",
Path::new("/nonexistent/my/tree"),
Duration::from_millis(60),
Duration::from_millis(20),
)
},
);
assert_eq!(confirmed, None, "a daemon error confirms nothing");
}
}