trusty-mpm 1.8.2

trusty-mpm: unified multi-agent orchestration platform (core, daemon, CLI, TUI, Telegram)
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
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
//! Fail-open check for the daemon's one credential read (#8236 item 7).
//!
//! Every store here is a fake and every credential is a literal written three
//! lines above the assertion that reads it. No test touches an OS keychain, a
//! launchd domain, a real `$HOME`, or the host's `.env.local`.
//!
//! What each test pins: the arm returns `None` AND says so at ERROR with the
//! variable name and the error kind. A downgrade of any arm — to a default
//! value, to a silent `None`, or to a different kind — fails here.

use std::sync::Mutex;
use std::sync::atomic::{AtomicUsize, Ordering};

use serial_test::serial;
use tracing::field::{Field, Visit};
use tracing_subscriber::layer::{Context, Layer};
use tracing_subscriber::prelude::*;
use trusty_common::credentials::KeyStoreError;

use super::test_env::EnvVarGuard;
use super::*;

/// Collected ERROR/WARN/INFO lines from one `with_default` scope.
#[derive(Default)]
struct Sink {
    /// One rendered line per event: level, then `name=value` per field.
    lines: Mutex<Vec<String>>,
}

/// A `tracing` layer that renders every event into [`Sink`].
struct CaptureLayer {
    /// Where rendered lines land.
    sink: Arc<Sink>,
}

/// Renders an event's fields into a flat string.
struct Fields {
    /// The line built so far.
    out: String,
}

impl Visit for Fields {
    fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
        use std::fmt::Write as _;
        let _ = write!(self.out, " {}={value:?}", field.name());
    }

    fn record_str(&mut self, field: &Field, value: &str) {
        use std::fmt::Write as _;
        let _ = write!(self.out, " {}={value}", field.name());
    }
}

impl<S: tracing::Subscriber> Layer<S> for CaptureLayer {
    fn on_event(&self, event: &tracing::Event<'_>, _ctx: Context<'_, S>) {
        let mut fields = Fields {
            out: event.metadata().level().to_string(),
        };
        event.record(&mut fields);
        self.sink
            .lines
            .lock()
            .unwrap_or_else(std::sync::PoisonError::into_inner)
            .push(fields.out);
    }
}

/// Run `f` with a capturing subscriber installed on this thread.
fn capture<T>(f: impl FnOnce() -> T) -> (T, Vec<String>) {
    let sink = Arc::new(Sink::default());
    let subscriber = tracing_subscriber::registry().with(CaptureLayer {
        sink: Arc::clone(&sink),
    });
    let out = tracing::subscriber::with_default(subscriber, f);
    let lines = sink
        .lines
        .lock()
        .unwrap_or_else(std::sync::PoisonError::into_inner)
        .clone();
    (out, lines)
}

/// Assert exactly one ERROR line, naming `var` and `kind`.
fn assert_logged(lines: &[String], var: &str, kind: &str) {
    let errors: Vec<&String> = lines.iter().filter(|l| l.starts_with("ERROR")).collect();
    assert_eq!(errors.len(), 1, "expected one ERROR line, got {lines:?}");
    assert!(errors[0].contains(var), "{}", errors[0]);
    assert!(
        errors[0].contains(&format!("kind={kind}")),
        "expected kind={kind}: {}",
        errors[0]
    );
}

/// A store that answers instantly and holds nothing, counting its reads.
struct AbsentStore {
    /// How many reads reached the store.
    calls: Arc<AtomicUsize>,
}

impl KeyStore for AbsentStore {
    fn get(&self, _provider: &str) -> Option<String> {
        self.calls.fetch_add(1, Ordering::SeqCst);
        None
    }
    fn set(&self, _provider: &str, _value: &str) -> Result<(), KeyStoreError> {
        Ok(())
    }
    fn unset(&self, _provider: &str) -> Result<(), KeyStoreError> {
        Ok(())
    }
    fn list(&self) -> Vec<String> {
        Vec::new()
    }
}

/// A store whose read never returns — a SecurityAgent dialog, in a fake.
struct NeverReturns;

impl KeyStore for NeverReturns {
    fn get(&self, _provider: &str) -> Option<String> {
        loop {
            std::thread::sleep(Duration::from_secs(3600));
        }
    }
    fn set(&self, _provider: &str, _value: &str) -> Result<(), KeyStoreError> {
        Ok(())
    }
    fn unset(&self, _provider: &str) -> Result<(), KeyStoreError> {
        Ok(())
    }
    fn list(&self) -> Vec<String> {
        Vec::new()
    }
}

/// A store whose backend refuses — a locked or denied keychain.
struct RefusingStore;

impl KeyStore for RefusingStore {
    fn get(&self, _provider: &str) -> Option<String> {
        None
    }
    fn try_get(&self, _provider: &str) -> Result<Option<String>, KeyStoreError> {
        Err(KeyStoreError::Keyring("locked".to_string()))
    }
    fn set(&self, _provider: &str, _value: &str) -> Result<(), KeyStoreError> {
        Ok(())
    }
    fn unset(&self, _provider: &str) -> Result<(), KeyStoreError> {
        Ok(())
    }
    fn list(&self) -> Vec<String> {
        Vec::new()
    }
}

/// Why: "nothing is configured" is the one failure an operator is allowed to
/// ignore, and it is also the one most easily turned into a silent `None` by a
/// later refactor. The ERROR line is what tells a stopped feature apart from a
/// feature nobody turned on; asserting on it means a downgrade to a quiet
/// return — or to any non-`None` default — fails here.
/// Test: this test.
#[test]
#[serial]
fn an_absent_secret_is_none_and_logs_absent() {
    let _guard = EnvVarGuard::unset("BITBUCKET_APP_PASSWORD");
    let calls = Arc::new(AtomicUsize::new(0));
    let store = Arc::new(AbsentStore {
        calls: Arc::clone(&calls),
    });

    let (resolved, lines) = capture(|| {
        resolve_secret_with("BITBUCKET_APP_PASSWORD", store, Duration::from_millis(500))
    });

    assert!(resolved.is_none(), "an absent credential must not resolve");
    assert_eq!(
        calls.load(Ordering::SeqCst),
        1,
        "the store was consulted once"
    );
    assert_logged(&lines, "BITBUCKET_APP_PASSWORD", "absent");
}

/// Why (#8236 item 6): under launchd the first read by a rebuilt binary parks
/// on a dialog. The daemon must come back `None` inside its bound and must say
/// TIMEOUT, not ABSENT — an operator told "absent" reinstates the plaintext
/// plist entry, which is this issue happening again. A downgrade of this arm to
/// a retry, a cached value, or the absent wording fails here.
/// Test: this test.
#[test]
#[serial]
fn a_store_timeout_is_none_and_logs_timeout() {
    let _guard = EnvVarGuard::unset("JIRA_API_TOKEN");
    let bound = Duration::from_millis(200);

    let started = std::time::Instant::now();
    let (resolved, lines) =
        capture(|| resolve_secret_with("JIRA_API_TOKEN", Arc::new(NeverReturns), bound));
    let waited = started.elapsed();

    assert!(resolved.is_none(), "a timed-out read must not resolve");
    assert!(
        waited < Duration::from_secs(2),
        "the caller waited {waited:?}"
    );
    assert_logged(&lines, "JIRA_API_TOKEN", "timeout");
}

/// Why: a refused store and an unconfigured one send an operator to two
/// different remedies. The kind is the only thing that separates them, so an
/// arm that collapsed every backend failure into "absent" — or, worse, into a
/// success with an empty value — fails here.
/// Test: this test.
#[test]
#[serial]
fn a_store_error_is_none_and_logs_the_kind() {
    let _guard = EnvVarGuard::unset("LINEAR_API_KEY");

    let (resolved, lines) = capture(|| {
        resolve_secret_with(
            "LINEAR_API_KEY",
            Arc::new(RefusingStore),
            Duration::from_millis(500),
        )
    });

    assert!(resolved.is_none(), "a refused store must not resolve");
    assert_logged(&lines, "LINEAR_API_KEY", "keyring-backend");
}

/// Why: `[llm] api_key_env` lets an operator name any variable, and an
/// unregistered one has no provider — so it has no store tier and no `.env`
/// tier either. Reinstating a file tier here is how a credential gets back into
/// a plaintext committed file, which is the same defect in a different file.
/// This pins that the store is NEVER consulted for such a name, and that the
/// process environment is the whole answer.
/// Test: this test.
#[test]
#[serial]
fn an_unregistered_variable_reads_only_the_process_environment() {
    let calls = Arc::new(AtomicUsize::new(0));
    let store = || {
        Arc::new(AbsentStore {
            calls: Arc::clone(&calls),
        })
    };

    let present = EnvVarGuard::set("TM_TEST_UNREGISTERED_8236", "from-the-process-env");
    let (resolved, lines) = capture(|| {
        resolve_secret_with(
            "TM_TEST_UNREGISTERED_8236",
            store(),
            Duration::from_millis(500),
        )
    });
    assert_eq!(resolved.as_deref(), Some("from-the-process-env"));
    assert!(lines.is_empty(), "a success must log nothing: {lines:?}");
    drop(present);

    let (resolved, lines) = capture(|| {
        resolve_secret_with(
            "TM_TEST_UNREGISTERED_8236",
            store(),
            Duration::from_millis(500),
        )
    });
    assert!(resolved.is_none(), "an unset variable must not resolve");
    assert_logged(&lines, "TM_TEST_UNREGISTERED_8236", "absent");
    assert_eq!(
        calls.load(Ordering::SeqCst),
        0,
        "an unregistered name must never reach the credential store"
    );
}

/// A store that holds a value for every provider, counting its reads.
struct HoldingStore {
    /// How many reads reached the store.
    calls: Arc<AtomicUsize>,
}

impl KeyStore for HoldingStore {
    fn get(&self, _provider: &str) -> Option<String> {
        self.calls.fetch_add(1, Ordering::SeqCst);
        Some("from-the-store-9121".to_string())
    }
    fn set(&self, _provider: &str, _value: &str) -> Result<(), KeyStoreError> {
        Ok(())
    }
    fn unset(&self, _provider: &str) -> Result<(), KeyStoreError> {
        Ok(())
    }
    fn list(&self) -> Vec<String> {
        Vec::new()
    }
}

/// The registered credential the two sandbox tests read.
///
/// Why: no other test in this crate touches it. A variable a non-`#[serial]`
/// test also mutates (`TELEGRAM_BOT_TOKEN` is one) races the guard below and
/// can send that test to the operator's real store.
const SANDBOX_TEST_VAR: &str = "FIREWORKS_API_KEY";

/// Why (#9121): a sandbox daemon with its own `$HOME` still reached the
/// operator's bot token, because the Keychain is per-user. In sandbox mode the
/// store must not be read at all — not read and then discarded.
/// What: a store holding a value; the gated read answers `None`, the store saw
/// zero reads, and the WARN names the variable without the value. No assertion
/// message prints a resolved value.
/// Test: this test.
#[test]
#[serial]
fn sandbox_mode_never_consults_the_store() {
    let _env = EnvVarGuard::unset(SANDBOX_TEST_VAR);
    let calls = Arc::new(AtomicUsize::new(0));
    let store = Arc::new(HoldingStore {
        calls: Arc::clone(&calls),
    });

    let (resolved, lines) = capture(|| {
        resolve_gated(true, SANDBOX_TEST_VAR, || {
            resolve_secret_with(SANDBOX_TEST_VAR, store, Duration::from_millis(500))
        })
    });

    assert!(resolved.is_none(), "a sandbox must resolve no credential");
    assert_eq!(calls.load(Ordering::SeqCst), 0, "the store was read");
    assert!(
        lines
            .iter()
            .any(|l| l.starts_with("WARN") && l.contains(SANDBOX_TEST_VAR)),
        "the skip must be logged by name: {lines:?}"
    );
    assert!(
        !lines.iter().any(|l| l.contains("from-the-store")),
        "no value may be logged: {lines:?}"
    );
}

/// Why: the sandbox test above passes trivially if the fake store were never
/// wired in. This is the same read with the gate open: the store IS consulted.
/// Test: this test.
#[test]
#[serial]
fn outside_sandbox_mode_the_store_is_consulted() {
    let _env = EnvVarGuard::unset(SANDBOX_TEST_VAR);
    let calls = Arc::new(AtomicUsize::new(0));
    let store = Arc::new(HoldingStore {
        calls: Arc::clone(&calls),
    });

    let resolved = resolve_gated(false, SANDBOX_TEST_VAR, || {
        resolve_secret_with(SANDBOX_TEST_VAR, store, Duration::from_millis(500))
    });

    assert!(
        resolved.as_deref() == Some("from-the-store-9121"),
        "the open gate did not return the store's value"
    );
    assert_eq!(calls.load(Ordering::SeqCst), 1);
}

/// Why (#9121): the classifier and the manager built their `Configurator`
/// store with `default_store()`, past the latch. In sandbox mode the real
/// store must not even be constructed — `default_store()` probes the Keychain.
/// What: the real constructor panics; the gated store is empty.
/// Test: this test.
#[test]
fn a_sandbox_store_never_builds_the_real_store() {
    let store = store_gated(true, || panic!("the real store was built in sandbox mode"));

    assert!(store.get("telegram").is_none());
    assert!(store.get("openrouter").is_none());
    assert!(store.list().is_empty(), "the sandbox store is not empty");
}

/// Why: the test above passes trivially if the gate always returned an empty
/// store. With the gate open, the real store IS what the caller gets.
/// Test: this test.
#[test]
fn outside_sandbox_mode_the_real_store_is_built() {
    let calls = Arc::new(AtomicUsize::new(0));
    let reads = Arc::clone(&calls);

    let store = store_gated(false, move || Box::new(HoldingStore { calls: reads }));

    assert!(
        store.get("openrouter").as_deref() == Some("from-the-store-9121"),
        "the open gate did not hand back the real store"
    );
    assert_eq!(calls.load(Ordering::SeqCst), 1);
}

/// Why (#9121): the `/activity` key-presence probe read
/// `resolve_env_var_bounded` directly, which walks `.env.local` and the
/// Keychain. Gated, it answers `Absent`, reads no tier, and — like the probe it
/// replaces (#8563) — logs nothing.
/// Test: this test.
#[test]
fn a_sandboxed_bounded_resolve_is_absent_and_reads_nothing() {
    let (outcome, lines) = capture(|| {
        bounded_gated(true, "OPENROUTER_API_KEY", |_| {
            panic!("a tier was read in sandbox mode")
        })
    });

    assert!(
        matches!(&outcome, Err(SecretResolveError::Absent { var }) if var == "OPENROUTER_API_KEY"),
        "a sandboxed probe must be Absent"
    );
    assert!(lines.is_empty(), "the probe must not log: {lines:?}");
}

/// Test: this test.
#[test]
fn outside_sandbox_mode_the_bounded_resolver_runs() {
    let outcome = bounded_gated(false, "OPENROUTER_API_KEY", |_| Ok("resolved".into()));

    assert!(
        outcome.as_deref() == Ok("resolved"),
        "the open gate did not run the resolver"
    );
}