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,
&WallClock(Instant::now()),
)
}
trait PollClock {
fn elapsed(&self) -> Duration;
fn sleep(&self, gap: Duration);
}
struct WallClock(Instant);
impl PollClock for WallClock {
fn elapsed(&self) -> Duration {
self.0.elapsed()
}
fn sleep(&self, gap: Duration) {
std::thread::sleep(gap);
}
}
fn confirm_within(
socket: &Path,
index_id: &str,
root: &Path,
deadline: Duration,
gap: Duration,
clock: &impl PollClock,
) -> Option<String> {
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(),
clock.elapsed()
);
return Some(registered);
}
let elapsed = clock.elapsed();
if elapsed >= deadline {
break;
}
clock.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::cell::{Cell, RefCell};
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() }]
})
}
fn wall() -> WallClock {
WallClock(Instant::now())
}
struct SteppedClock {
deadline: Duration,
now: Cell<Duration>,
sleeps: RefCell<Vec<Duration>>,
}
impl PollClock for SteppedClock {
fn elapsed(&self) -> Duration {
self.now.get()
}
fn sleep(&self, gap: Duration) {
let now = self.now.get();
assert!(
now < self.deadline,
"slept at {now:?}, at or past the deadline"
);
assert!(
now + gap <= self.deadline,
"a {gap:?} sleep at {now:?} outlives the {:?} deadline",
self.deadline
);
self.now.set(now + gap);
self.sleeps.borrow_mut().push(gap);
}
}
#[test]
fn confirm_within_stops_at_the_deadline_and_withholds() {
let reads = Arc::new(AtomicU32::new(0));
let counter = Arc::clone(&reads);
let deadline = Duration::from_millis(100);
let clock = SteppedClock {
deadline,
now: Cell::new(Duration::ZERO),
sleeps: RefCell::new(Vec::new()),
};
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"),
deadline,
Duration::from_millis(30),
&clock,
)
},
);
assert_eq!(
confirmed, None,
"an index the daemon never registers must stay unpinned (#5091)"
);
let ms = Duration::from_millis;
assert_eq!(
*clock.sleeps.borrow(),
[ms(30), ms(30), ms(30), ms(10)],
"the poll must sleep full gaps, clamp the last one, and stop at its deadline"
);
assert_eq!(
reads.load(Ordering::SeqCst),
5,
"the confirm must read the registry once per pass, before and after every sleep"
);
}
#[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),
&wall(),
)
})
},
);
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),
&wall(),
)
},
);
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),
&wall(),
)
},
);
assert_eq!(confirmed, None, "a daemon error confirms nothing");
}
}