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
//! M4b E2E (spec §11.3): a REAL in-process mesh session lands audit records with the arguments
//! HASHED (the raw argument bytes never appear in the file) and the session + proxied-request event
//! classes present. The trust class is proven by `trust_mutations_emit_audit_events` and the
//! blob-fetch class by `served_get_records_blob_fetch_audit` (both Task 5).
//!
//! Harness: a direct `serve(...)` + `mcpmesh_net::connect(...)` mesh over two localhost endpoints (no
//! subprocess, no control socket). `connect` dials the FULL server address, so there is no id→addr
//! discovery step and no `MemoryLookup`. The audit sink writes to a hermetic temp dir.
use std::sync::Arc;
use std::time::Duration;
use mcpmesh::allowlist::{AllowlistGate, PeerEntry, PeerStore};
use mcpmesh::audit::{AuditLog, AuditSink};
use mcpmesh::config::Config;
use mcpmesh::daemon::build_services_audited;
use mcpmesh_net::{TrustGate, serve};
use serde_json::json;
use tokio::time::timeout;
const STUB: &str = env!("CARGO_BIN_EXE_echo_mcp_stub");
async fn local_endpoint() -> iroh::Endpoint {
iroh::Endpoint::builder(iroh::endpoint::presets::Minimal)
.relay_mode(iroh::RelayMode::Disabled)
.alpns(vec![mcpmesh_net::ALPN_MCP.to_vec()])
.bind()
.await
.expect("bind localhost endpoint")
}
#[tokio::test]
async fn real_session_audits_with_hashed_args_and_all_event_classes() {
timeout(Duration::from_secs(90), async {
let secret = "TOP-SECRET-ARGUMENT-do-not-log-me";
let dir = tempfile::tempdir().unwrap();
let audit_dir = dir.path().join("audit");
// The serving "server machine" endpoint: an echo `run` service audited to `audit_dir`.
// Capture the address BEFORE `server_ep` is moved into `serve` below.
let server_ep = local_endpoint().await;
let server_addr = server_ep.addr();
// Trust the caller endpoint (the connecting side) as nickname "bob". `PeerEntry` is a struct
// literal (allowlist.rs) — fields are `endpoint_id` / `nickname` / `services` / `paired_at`
// (mirrors proxy_roundtrip.rs); there is no `PeerEntry::new`.
let caller_ep = local_endpoint().await;
let caller_id = *caller_ep.id().as_bytes();
let gate: Arc<dyn TrustGate> = Arc::new(AllowlistGate::new({
let store = Arc::new(PeerStore::open(&dir.path().join("state.redb")).unwrap());
store
.add(PeerEntry {
endpoint_id: caller_id,
nickname: "bob".into(),
services: vec!["notes".into()],
paired_at: None,
user_id: None,
last_addr: None,
})
.unwrap();
store
}));
// Serve an echo `run` service named "notes", threaded with a REAL audit sink.
// TOML LITERAL string for the stub path (like every sibling suite): a windows
// `D:\…\stub.exe` path in a basic "…" string is an invalid escape sequence.
// 0.8.0 (#38): `allow` holds STABLE principals — the caller's `eid:<hex>` — never the
// display nickname ("bob" is display-only and would not admit).
let caller_eid = format!("eid:{}", caller_ep.id());
let cfg = Config::from_toml_str(&format!(
"[services.notes]\nrun = ['{STUB}']\nallow = [\"{caller_eid}\"]\n"
))
.unwrap();
let sink = AuditSink::new(AuditLog::spawn(audit_dir.clone()));
let _serve = serve(
server_ep,
gate,
build_services_audited(&cfg, &sink, &mcpmesh::limits::MeshLimiters::unlimited()),
Arc::new(mcpmesh_net::ConnRegistry::new()),
);
// Drive one MCP session over the mesh from the caller endpoint: connect (dialing the FULL
// server address — no id→addr discovery, no MemoryLookup), initialize, then a tools/call whose
// arguments carry `secret`.
let mut transport = mcpmesh_net::connect(&caller_ep, server_addr, "notes")
.await
.unwrap()
.0;
// initialize
transport
.send_value(json!({
"jsonrpc":"2.0","id":1,"method":"initialize",
"params":{"_meta":{"mcpmesh/service":"notes"},"capabilities":{}}
}))
.await
.unwrap();
let _init = transport.recv_value().await.unwrap().unwrap();
// tools/call with the sensitive argument.
transport
.send_value(json!({
"jsonrpc":"2.0","id":2,"method":"tools/call",
"params":{"name":"read_file","arguments":{"text":secret}}
}))
.await
.unwrap();
let call = transport.recv_value().await.unwrap().unwrap();
assert_eq!(call["id"], 2);
drop(transport); // EOF → session_close
// Poll the audit file until the request record lands.
let month = &mcpmesh::audit::now_ts()[..7];
let file = audit_dir.join(format!("{month}.jsonl"));
let mut body = String::new();
for _ in 0..100 {
if let Ok(b) = std::fs::read_to_string(&file)
&& b.contains("\"kind\":\"request\"")
&& b.contains("\"kind\":\"session_open\"")
&& b.contains("\"kind\":\"session_close\"")
{
body = b;
break;
}
tokio::time::sleep(Duration::from_millis(30)).await;
}
// PRIVACY (the headline): the raw argument NEVER appears anywhere in the audit file.
assert!(
!body.contains(secret),
"raw argument leaked into the audit log:\n{body}"
);
// The request record carries the tool NAME + a blake3 args hash + a bytes_out count.
assert!(body.contains("\"tool\":\"read_file\""));
assert!(body.contains("\"method\":\"tools/call\""));
assert!(body.contains("blake3:"));
assert!(body.contains("\"bytes_out\":"));
assert!(body.contains("\"peer\":\"bob\""));
// Session lifecycle present.
assert!(body.contains("\"kind\":\"session_open\""));
assert!(body.contains("\"kind\":\"session_close\""));
// #57: EVERY line of this session carries the caller's STABLE device principal
// alongside the display name — the same eid: the gate resolved (session/request/blob
// records key the DEVICE; a bound peer's b64u: appears on the trust record and joins
// via the status peers list). open + request + close must all agree (the close reads
// the stored session row), so the count pins all three producers at once.
assert_eq!(
body.matches(&format!("\"principal\":\"{caller_eid}\""))
.count(),
3,
"session_open + request + session_close each carry the caller's eid (#57): {body}"
);
})
.await
.expect("audit E2E timed out");
}