1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
//! Unit tests for the shared trusty-search socket client (#7237).
//!
//! Why: isolated in a sibling file (declared via `#[path =
//! "search_rpc_tests.rs"] mod tests;` in `search_rpc.rs`) so the production
//! module stays well under the 500-SLOC cap. As a child module, `super::`
//! reaches its private items.
//!
//! What: the socket-path override, the dial-failure arm both entry points
//! share, and the blocking wrapper's two answering arms — a `result` and the
//! daemon's own coded refusal.
//!
//! Test: `cargo test -p trusty-common --features uds -- search_rpc::tests`
use super::*;
use crate::uds_mock::{self, RpcError};
/// Long enough that a loaded machine's loopback round trip is not a flake,
/// short enough that a dead socket does not stall the suite.
const PROBE_TIMEOUT: Duration = Duration::from_secs(5);
/// Why: the override is the only way a probe can be pointed at a daemon a
/// test started without redirecting every other trusty-* client in the
/// process through `TRUSTY_DATA_DIR_OVERRIDE`.
///
/// It holds `crate::data_dir::ENV_LOCK` rather than `#[serial_test::serial]`,
/// even though nothing here reads the data dir: `search_index`'s rigs set this
/// same variable under that lock, and two DIFFERENT mutexes serialise nothing
/// against each other — this test's `remove_var` wiped a running rig's override
/// mid-registration and it reported `DaemonUnreachable`.
/// Test: itself.
#[test]
fn search_socket_honours_the_env_override() {
let _guard = crate::data_dir::ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
// SAFETY: guarded by ENV_LOCK; removed below before returning.
unsafe {
std::env::set_var(TRUSTY_SEARCH_SOCKET_ENV, "/tmp/example-search.sock");
}
let resolved = search_socket();
unsafe {
std::env::remove_var(TRUSTY_SEARCH_SOCKET_ENV);
}
assert_eq!(
resolved.expect("an override always resolves"),
PathBuf::from("/tmp/example-search.sock")
);
}
/// Why: #9214 — the daemon binds `$TRUSTY_DATA_DIR/trusty-search.sock`
/// (`trusty_search::service::socket::resolve_socket_path`), so a client that
/// ignored the variable dialled the shared socket and could never reach an
/// isolated daemon it ran beside.
///
/// It sets the variable under `crate::data_dir::ENV_LOCK`, the lock every
/// env-mutating test in this crate shares, and restores both variables.
/// Test: itself.
#[test]
fn search_socket_follows_trusty_data_dir_like_the_daemon() {
let _guard = crate::data_dir::ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let tmp = tempfile::tempdir().expect("tempdir");
let saved_dir = std::env::var_os("TRUSTY_DATA_DIR");
let saved_socket = std::env::var_os(TRUSTY_SEARCH_SOCKET_ENV);
// SAFETY: guarded by ENV_LOCK; both variables are restored below.
unsafe {
std::env::set_var("TRUSTY_DATA_DIR", tmp.path());
std::env::remove_var(TRUSTY_SEARCH_SOCKET_ENV);
}
let resolved = search_socket();
unsafe {
match saved_dir {
Some(v) => std::env::set_var("TRUSTY_DATA_DIR", v),
None => std::env::remove_var("TRUSTY_DATA_DIR"),
}
if let Some(v) = saved_socket {
std::env::set_var(TRUSTY_SEARCH_SOCKET_ENV, v);
}
}
assert_eq!(
resolved.expect("an absolute TRUSTY_DATA_DIR resolves"),
tmp.path().join("trusty-search.sock"),
"the client must dial the socket the daemon binds under TRUSTY_DATA_DIR"
);
}
/// Why: #9214 — the shared rule's two edge arms. `var_os` reports an
/// exported-but-empty `TRUSTY_DATA_DIR` as `Some("")`, which must fall back to
/// the shared socket; a relative value would resolve against the cwd, so it is
/// refused. Holds `ENV_LOCK` because the shared arm reads
/// `TRUSTY_DATA_DIR_OVERRIDE`, which other tests set under that lock.
/// Test: itself.
#[test]
fn search_socket_under_treats_empty_as_unset_and_refuses_relative() {
let _guard = crate::data_dir::ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let shared = search_socket_under(None).expect("the shared path resolves");
let empty = std::ffi::OsString::new();
assert_eq!(
search_socket_under(Some(empty.as_os_str())).expect("empty is not fatal"),
shared,
"an empty override must be treated as unset"
);
let relative = std::ffi::OsString::from("relative/data-dir");
let err = search_socket_under(Some(relative.as_os_str()))
.expect_err("a relative override must be refused");
assert!(err.to_string().contains("absolute"), "{err:#}");
}
/// Why: every probe's failure arm depends on a dead socket FAILING rather
/// than consuming the budget — `tm doctor` must not hang on an absent
/// daemon, and `search_index`'s registration must not stall a session launch.
/// Test: itself.
#[tokio::test]
async fn call_at_reports_a_dead_socket_rather_than_hanging() {
let dir = tempfile::tempdir().expect("tempdir");
let socket = dir.path().join("absent.sock");
let err = call_at(&socket, METHOD_HEALTH, serde_json::json!({}), PROBE_TIMEOUT)
.await
.expect_err("nothing is listening");
assert!(
err.downcast_ref::<SearchRpcError>().is_none(),
"a dial failure is a transport error, never a daemon refusal: {err:#}"
);
}
/// Why: [`call_blocking`] is the entry point `search_index` uses from a
/// synchronous function that is itself frequently called from inside a tokio
/// runtime. Driving it from `spawn_blocking` inside a `#[tokio::test]` is that
/// arrangement exactly — if the wrapper built its runtime on the caller's
/// thread instead of its own, this would panic rather than answer.
/// Test: itself.
#[tokio::test]
async fn call_blocking_round_trips_against_a_listening_daemon() {
let daemon = uds_mock::spawn(|method, params| {
let answer = serde_json::json!({ "method": method, "params": params });
Box::pin(async move { Ok(answer) })
})
.await;
let socket = daemon.socket().to_path_buf();
let value = tokio::task::spawn_blocking(move || {
call_blocking(
&socket,
METHOD_INDEX_CREATE,
serde_json::json!({ "id": "api" }),
PROBE_TIMEOUT,
)
})
.await
.expect("the blocking worker must not panic")
.expect("the daemon answered");
assert_eq!(
value.get("method").and_then(serde_json::Value::as_str),
Some(METHOD_INDEX_CREATE),
"the method name must reach the daemon verbatim: {value}"
);
assert_eq!(
value
.pointer("/params/id")
.and_then(serde_json::Value::as_str),
Some("api"),
"the params must reach the daemon verbatim: {value}"
);
}
/// Why: the caller that matters branches on the CODE — a `409` is recoverable
/// (#6864) and a `404` is a definite verdict (#5045). A refusal flattened into
/// a formatted string would make both indistinguishable from a transport
/// failure.
/// Test: itself.
#[tokio::test]
async fn call_blocking_carries_the_daemons_own_error_code() {
let daemon = uds_mock::spawn(|_method, _params| {
Box::pin(async move { Err(RpcError::new(CODE_CONFLICT, "root_path is taken")) })
})
.await;
let socket = daemon.socket().to_path_buf();
let err = tokio::task::spawn_blocking(move || {
call_blocking(
&socket,
METHOD_INDEX_CREATE,
serde_json::json!({}),
PROBE_TIMEOUT,
)
})
.await
.expect("the blocking worker must not panic")
.expect_err("the daemon refused");
let refusal = err
.downcast_ref::<SearchRpcError>()
.expect("a daemon refusal must arrive typed, not as an opaque string");
assert!(refusal.is_conflict(), "code was {}", refusal.code);
assert!(!refusal.is_not_found());
assert_eq!(refusal.message, "root_path is taken");
}
/// Why: #6285 moves the MCP bridge onto this client, and its INDEX_UNAVAILABLE
/// contract relays the daemon's structured refusal body. A `data` member the
/// daemon sent and `call_at` dropped would strand every field but the code.
/// Test: itself.
#[tokio::test]
async fn call_at_carries_the_daemons_error_data() {
let body = serde_json::json!({
"error": "index_not_resident",
"index_id": "wt-1",
"retryable": true,
"restore_via": "search.query",
});
let sent = body.clone();
let daemon = uds_mock::spawn(move |_method, _params| {
let refusal = RpcError::new(-32002, "index_not_resident").with_data(sent.clone());
Box::pin(async move { Err(refusal) })
})
.await;
let err = call_at(
daemon.socket(),
METHOD_HEALTH,
serde_json::json!({}),
PROBE_TIMEOUT,
)
.await
.expect_err("the daemon refused");
let refusal = err
.downcast_ref::<SearchRpcError>()
.expect("a daemon refusal must arrive typed");
assert_eq!(refusal.data.as_ref(), Some(&body));
}
/// Why: a handler that dies mid-call is the one failure that reaches
/// [`call_blocking`] as neither a coded refusal nor a clean dial failure — the
/// daemon accepted the connection and then never answered. The caller is a
/// best-effort registration on a session-launch path, so what it owes is an
/// `Err` that names the method it was making, not a stall and not a panic
/// escaping the worker thread (#7237).
/// Test: itself.
#[tokio::test]
async fn call_blocking_reports_a_panicking_handler_rather_than_hanging() {
let daemon = uds_mock::spawn(|_method, _params| {
Box::pin(async move { panic!("the mock daemon handler died before answering") })
})
.await;
let socket = daemon.socket().to_path_buf();
let err = tokio::task::spawn_blocking(move || {
call_blocking(
&socket,
METHOD_INDEX_CREATE,
serde_json::json!({}),
PROBE_TIMEOUT,
)
})
.await
.expect("the panic must stay inside the daemon's connection task")
.expect_err("a handler that never answers cannot produce a result");
assert!(
err.downcast_ref::<SearchRpcError>().is_none(),
"an unanswered call is a transport failure, never a daemon refusal: {err:#}"
);
assert!(
format!("{err:#}").contains(METHOD_INDEX_CREATE),
"the error must name the method that failed: {err:#}"
);
}
/// Why: an absent socket is what "trusty-search is not running" looks like, and
/// the registration path treats it as fail-closed rather than falling back to
/// anything. It must come back promptly and as a transport failure.
/// Test: itself.
#[test]
fn call_blocking_reports_a_dead_socket_rather_than_hanging() {
let dir = tempfile::tempdir().expect("tempdir");
let socket = dir.path().join("absent.sock");
let err = call_blocking(&socket, METHOD_HEALTH, serde_json::json!({}), PROBE_TIMEOUT)
.expect_err("nothing is listening");
assert!(
err.downcast_ref::<SearchRpcError>().is_none(),
"a dial failure is a transport error, never a daemon refusal: {err:#}"
);
}