trusty-common 0.53.0

Shared utilities and provider-agnostic streaming chat (ChatProvider, OllamaProvider, OpenRouter, tool-use) for trusty-* projects
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
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
//! Confirming a `search.index.create` the daemon never answered in time
//! (#7237).
//!
//! Why: [`super::CREATE_TIMEOUT`] caps the create at one second because it runs
//! on a session-launch hot path, and until now a create the daemon completed
//! LATE was indistinguishable from one it refused. On 2026-09-10 the `writing`
//! project's index had been cold-parked with 32,754 chunks. The create arrived
//! while the daemon was reloading it, the daemon registered the index 3.8 s
//! later, and by then the client had already recorded `NotConfirmed`, withheld
//! the id, and fired a reindex the daemon answered with `unknown index: writing`.
//! The session then ran its whole length unpinned, against an index that
//! existed.
//!
//! What: [`confirm_after_no_answer`] re-reads the daemon's registry until an
//! index whose `root_path` IS this tree appears, or [`CONFIRM_DEADLINE`]
//! elapses. It is reached ONLY from [`super::reconcile::CreateOutcome::Unanswered`]
//! — a call that timed out, that the peer hung up on, or that failed mid-read —
//! never from a refusal, so no daemon error can become a confirmed index. The
//! match is [`super::reconcile::index_id_serving_root`], the same root-identity
//! comparison the #6864 collision recovery uses, so a confirmation names an
//! index the daemon actually holds for this root rather than one it was merely
//! asked for.
//!
//! The trade-off this picks: a session launch pays roughly [`CONFIRM_DEADLINE`]
//! on top of the create's own budget — see that constant for why one in-flight
//! registry read can carry it a second further — and only when the daemon left
//! the create unanswered. An ordinary create answers in milliseconds and never reaches
//! here. A launch that exhausts the deadline is told so in one `warn` and still
//! withholds the pin — the daemon may finish registering afterwards, and nothing
//! here retries.
//!
//! Test: `confirm_within_stops_at_the_deadline_and_withholds`,
//! `confirm_within_never_confirms_an_index_at_another_root`,
//! `an_exhausted_confirm_deadline_is_still_a_warning` below, plus
//! `a_late_create_is_confirmed_by_polling_the_registry`,
//! `a_create_the_daemon_hung_up_on_is_confirmed_by_polling_the_registry` and
//! `an_unconfirmed_registration_fires_no_reindex` in `search_index_tests.rs`.

use std::path::Path;
use std::time::{Duration, Instant};

use super::reconcile::{ListFailure, fetch_index_list, index_id_serving_root};

/// How long [`confirm_after_no_answer`] keeps asking, before giving up (#7237).
///
/// Sized against the failure it exists for: the measured cold reload took
/// 3.8 s from the client's first byte, of which the create's own
/// [`super::CREATE_TIMEOUT`] covers the first second. Four seconds leaves
/// margin over the rest without turning "the daemon is wedged" into a launch
/// that hangs.
///
/// Not the total wall time: the deadline is tested only after a registry read
/// returns, and each read carries its own one-second budget, so a read that
/// starts just under the deadline can push the confirm to roughly five seconds.
/// The alternative — cancelling a read mid-flight — would throw away the answer
/// this poll exists to get.
const CONFIRM_DEADLINE: Duration = Duration::from_secs(4);

/// How long [`confirm_after_no_answer`] waits between two registry reads.
///
/// A read that FAILS has already spent its own ~1 s budget, so this gap only
/// paces the case where the daemon answers promptly with a registry that does
/// not name this tree yet.
const CONFIRM_POLL_GAP: Duration = Duration::from_millis(250);

/// The id the daemon ended up registering for `root`, or `None` (#7237).
///
/// Why: see the module doc. A create the daemon never answered is not evidence
/// that it refused, and the registry is the one place that settles which of the
/// two happened.
/// What: polls `search.indexes.list` through [`fetch_index_list`] until
/// [`index_id_serving_root`] names an index for `root`, then returns that id —
/// which may differ from `index_id`, exactly as the #6864 recovery's does. Gives
/// up at [`CONFIRM_DEADLINE`] with one `warn` and `None`, which leaves the
/// caller's pin unadvanced. Never propagates an error and never sleeps past the
/// deadline.
/// Test: `confirm_within_stops_at_the_deadline_and_withholds`,
/// `a_late_create_is_confirmed_by_polling_the_registry`,
/// `a_create_the_daemon_hung_up_on_is_confirmed_by_polling_the_registry`.
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()),
    )
}

/// The time source [`confirm_within`] measures its deadline against.
///
/// Why: #8284 — on the wall clock, one slow registry read under CI load spent
/// the whole test deadline, so the poll loop's own schedule could not be
/// asserted. A test supplies a clock that moves only when the loop sleeps.
/// What: `elapsed` is time since the confirm started; `sleep` waits `gap`.
/// Test: `confirm_within_stops_at_the_deadline_and_withholds`.
trait PollClock {
    fn elapsed(&self) -> Duration;
    fn sleep(&self, gap: Duration);
}

/// The production [`PollClock`]: [`Instant`] plus [`std::thread::sleep`].
struct WallClock(Instant);

impl PollClock for WallClock {
    fn elapsed(&self) -> Duration {
        self.0.elapsed()
    }

    fn sleep(&self, gap: Duration) {
        std::thread::sleep(gap);
    }
}

/// [`confirm_after_no_answer`] with the two budgets and the clock supplied.
///
/// Why: the real deadline is seconds long, so a test that drove it would spend
/// them. Taking both as parameters lets the give-up and wrong-root arms be
/// asserted in milliseconds while the one end-to-end test still exercises the
/// production constants.
/// What: read, match, sleep, repeat — the sleep is clamped to whatever is left
/// of `deadline` so no wait is started that outlives it, and the loop always
/// performs at least one read even when `deadline` is zero. `deadline` is tested
/// between reads, not during one, so the final read can still overrun it by its
/// own budget; see [`CONFIRM_DEADLINE`]. Time is read from `clock`.
/// Test: `confirm_within_stops_at_the_deadline_and_withholds`,
/// `confirm_within_never_confirms_an_index_at_another_root`,
/// `an_exhausted_confirm_deadline_is_still_a_warning`.
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;
        // #7237: matched on the TREE, so only an index the daemon really holds
        // for this root can confirm the registration.
        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};

    /// Run `body` against a mock daemon answering every method through
    /// `handler`, on a socket this function owns.
    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
    }

    /// An `indexes` listing naming one entry.
    fn listing(id: &str, root: &Path) -> serde_json::Value {
        serde_json::json!({
            "indexes": [{ "id": id, "root_path": root.to_string_lossy() }]
        })
    }

    /// The wall clock, for tests that assert no timing.
    fn wall() -> WallClock {
        WallClock(Instant::now())
    }

    /// A [`PollClock`] that advances only when the loop sleeps (#8284).
    ///
    /// Why: a registry read's wall time is the host's, not the loop's. Frozen
    /// across reads, the clock makes the poll schedule a function of `deadline`
    /// and `gap` alone.
    /// What: records each sleep; panics on a sleep that starts at or ends past
    /// `deadline`, so a loop that ignores its deadline fails instead of hanging.
    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);
        }
    }

    /// A daemon that never registers this tree leaves the pin unadvanced
    /// (#7237).
    ///
    /// Why: requirement 3 of the fix — the withholding #5091 introduced must
    /// survive it. A create that went unanswered because the daemon is wedged,
    /// or because the id is genuinely never going to exist, must still end in
    /// `None`, and the poll must stop rather than run forever.
    /// What: a daemon whose registry stays empty, on a [`SteppedClock`] with a
    /// 100 ms deadline and a 30 ms gap; asserts `None`, the exact sleep schedule
    /// (three full gaps, then one clamped to the 10 ms left), and one daemon read
    /// per loop pass.
    /// Test: this test.
    #[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);
        // #8284: a stepped clock, not the wall clock — one slow read under CI
        // load used to spend the whole deadline and leave a single read.
        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"
        );
    }

    /// A confirm that exhausts its deadline is still a warning (#7390).
    ///
    /// Why: #7390 lowered the line that ANNOUNCES the registry check, on the
    /// grounds that it fires before the outcome is known. The give-up here is
    /// the outcome — the id is withheld, the session runs unpinned, and this
    /// warn is the only signal an operator gets — so a fix that quieted it too
    /// would trade one false alarm for a silent failure. Pinned beside the
    /// zero-warning assertion in
    /// `super::super::reconcile::tests::an_unanswered_create_the_registry_confirms_emits_no_warning`,
    /// which is what would push a well-meant edit into this line.
    /// What: a daemon whose registry stays empty, under a capturing subscriber
    /// and a millisecond deadline; asserts the withheld id still comes with a
    /// WARN line naming the tree.
    /// Test: itself.
    #[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]
        );
    }

    /// A registry entry for ANOTHER tree never confirms this one (#7237).
    ///
    /// Why: the fail-open check. The whole confirm exists to turn a silent
    /// daemon into a pinnable id, and the one way that could go wrong is
    /// accepting some other index as evidence. The match is on the tree, so a
    /// busy daemon full of other projects confirms nothing here.
    /// What: a daemon whose registry names one index at an unrelated root;
    /// asserts `None`.
    /// Test: this test.
    #[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"
        );
    }

    /// A refusing registry is a failed read, not a confirmation (#7237).
    ///
    /// Why: the second half of the fail-open check — a daemon ERROR must never
    /// read as "the index is there". [`fetch_index_list`] answers `None` for a
    /// refusal, and the loop has to treat that as "not yet", not as a match.
    /// What: a daemon that refuses every method; asserts `None`.
    /// Test: this test.
    #[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");
    }
}