use std::path::{Path, PathBuf};
use crate::connector::{ServiceConnector, ServiceInfo, ServiceLifecycle, ServiceStatus};
use crate::search_uds::{HEALTH_TIMEOUT, METHOD_HEALTH, SEARCH_SERVICE};
use super::helpers::binary_on_path;
#[derive(Debug, serde::Deserialize)]
struct HealthEnvelope {
#[allow(dead_code)]
status: String,
version: Option<String>,
}
pub struct SearchConnector {
socket: Option<PathBuf>,
}
impl SearchConnector {
pub fn new() -> Self {
Self { socket: None }
}
pub fn with_socket(socket: PathBuf) -> Self {
Self {
socket: Some(socket),
}
}
fn socket_path(&self) -> Result<PathBuf, String> {
match &self.socket {
Some(p) => Ok(p.clone()),
None => crate::search_uds::socket_path(),
}
}
}
impl Default for SearchConnector {
fn default() -> Self {
Self::new()
}
}
fn probe_health(socket: &Path) -> Option<HealthEnvelope> {
let socket = socket.to_path_buf();
let handle = std::thread::Builder::new()
.name("console-search-probe".to_owned())
.spawn(move || {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.ok()?;
rt.block_on(async {
let result = crate::search_uds::call(
&socket,
METHOD_HEALTH,
serde_json::json!({}),
HEALTH_TIMEOUT,
)
.await
.ok()?;
serde_json::from_value::<HealthEnvelope>(result).ok()
})
})
.ok()?;
handle.join().ok()?
}
impl ServiceConnector for SearchConnector {
fn id(&self) -> &'static str {
SEARCH_SERVICE
}
fn display_name(&self) -> &'static str {
"Trusty Search"
}
fn detect(&self) -> ServiceInfo {
self.detect_from(self.socket_path())
}
}
impl SearchConnector {
fn detect_from(&self, socket: Result<PathBuf, String>) -> ServiceInfo {
let base =
|status: ServiceStatus, version: Option<String>, hint: Option<String>| ServiceInfo {
id: self.id().to_string(),
display_name: self.display_name().to_string(),
status,
version,
url: None,
hint,
lifecycle: ServiceLifecycle::Daemon,
};
if !binary_on_path(SEARCH_SERVICE) {
return base(ServiceStatus::Absent, None, None);
}
let socket = match socket {
Ok(p) => p,
Err(reason) => return base(ServiceStatus::Available, None, Some(reason)),
};
match probe_health(&socket) {
Some(health) => base(ServiceStatus::Running, health.version, None),
None => base(ServiceStatus::Available, None, None),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::path::PathBuf;
fn installed() -> ServiceStatus {
if which::which(SEARCH_SERVICE).is_ok() {
ServiceStatus::Available
} else {
ServiceStatus::Absent
}
}
fn stub_daemon(dir: &Path, reply: impl Into<String>) -> PathBuf {
let socket = dir.join("sockets").join("search.sock");
let reply = reply.into();
let listener = trusty_common::uds::bind_hardened(&socket).expect("bind");
tokio::spawn(async move {
use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
let Ok((mut conn, _)) = listener.accept().await else {
return;
};
let mut sink = Vec::new();
let _ = conn.read_to_end(&mut sink).await;
let _ = conn.write_all(reply.as_bytes()).await;
let _ = conn.write_all(b"\n").await;
let _ = conn.flush().await;
});
socket
}
async fn detect_against(socket: PathBuf) -> ServiceInfo {
tokio::task::spawn_blocking(move || SearchConnector::with_socket(socket).detect())
.await
.expect("detect")
}
#[test]
fn search_connector_reports_available_when_nothing_is_serving() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let connector = SearchConnector::with_socket(tmp.path().join("absent.sock"));
let info = connector.detect();
assert_eq!(info.status, installed());
assert_eq!(info.id, SEARCH_SERVICE);
assert_eq!(info.display_name, "Trusty Search");
assert!(
info.url.is_none(),
"a UDS daemon has no URL to link to: {info:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn search_connector_reads_the_version_off_a_live_socket() {
if which::which(SEARCH_SERVICE).is_err() {
return;
}
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_daemon(
tmp.path(),
r#"{"jsonrpc":"2.0","id":1,"result":{"status":"ok","version":"0.49.6","indexes":3}}"#,
);
let info = detect_against(socket).await;
assert_eq!(info.status, ServiceStatus::Running);
assert_eq!(info.version.as_deref(), Some("0.49.6"));
}
#[tokio::test(flavor = "multi_thread")]
async fn search_connector_reports_an_error_frame_as_not_running() {
if which::which(SEARCH_SERVICE).is_err() {
return;
}
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_daemon(
tmp.path(),
r#"{"jsonrpc":"2.0","id":1,"error":{"code":-32601,"message":"method not found"}}"#,
);
let info = detect_against(socket).await;
assert_eq!(info.status, ServiceStatus::Available);
assert!(info.version.is_none(), "{info:?}");
}
#[tokio::test(flavor = "multi_thread")]
async fn search_connector_reports_a_non_health_answer_as_not_running() {
if which::which(SEARCH_SERVICE).is_err() {
return;
}
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = stub_daemon(tmp.path(), r#"{"jsonrpc":"2.0","id":1,"result":{"hi":1}}"#);
let info = detect_against(socket).await;
assert_eq!(info.status, ServiceStatus::Available);
}
#[test]
fn search_connector_surfaces_an_unresolvable_socket_path_as_a_hint() {
let connector = SearchConnector::new();
let info = connector.detect_from(Err("no data directory".to_string()));
if info.status == ServiceStatus::Absent {
return;
}
assert_eq!(info.status, ServiceStatus::Available);
assert_eq!(info.hint.as_deref(), Some("no data directory"));
}
}