1use crate::contact_directory::{ContactEntry, MobTransport};
53
54#[derive(Debug, Clone, PartialEq, Eq)]
56pub enum RemoteMobError {
57 UnsupportedTransport { mob_id: String, transport: String },
59 ControlChannelUnavailable {
63 mob_id: String,
64 endpoint: String,
65 operation: &'static str,
66 },
67 Rejected {
71 mob_id: String,
72 endpoint: String,
73 code: String,
74 message: String,
75 },
76 Encode { endpoint: String, message: String },
79 Decode { endpoint: String, message: String },
82}
83
84impl RemoteMobError {
85 pub(crate) fn with_context(self, context: String) -> Self {
89 match self {
90 Self::ControlChannelUnavailable {
91 mob_id,
92 endpoint,
93 operation,
94 } => Self::ControlChannelUnavailable {
95 mob_id: if mob_id.is_empty() { context } else { mob_id },
96 endpoint,
97 operation,
98 },
99 other => other,
100 }
101 }
102}
103
104impl std::fmt::Display for RemoteMobError {
105 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
106 match self {
107 Self::UnsupportedTransport { mob_id, transport } => {
108 write!(
109 f,
110 "remote mob '{mob_id}' uses unsupported transport: {transport}"
111 )
112 }
113 Self::ControlChannelUnavailable {
114 mob_id,
115 endpoint,
116 operation,
117 } => {
118 write!(
119 f,
120 "remote mob '{mob_id}' control channel ({endpoint}): operation '{operation}' \
121 failed — confirm the peer gateway is running with a control listener bound \
122 on this endpoint"
123 )
124 }
125 Self::Rejected {
126 mob_id,
127 endpoint,
128 code,
129 message,
130 } => {
131 write!(
132 f,
133 "remote mob '{mob_id}' control channel ({endpoint}) rejected request \
134 [{code}]: {message}"
135 )
136 }
137 Self::Encode { endpoint, message } => {
138 write!(f, "control request encode failed for {endpoint}: {message}")
139 }
140 Self::Decode { endpoint, message } => {
141 write!(
142 f,
143 "control response decode failed for {endpoint}: {message}"
144 )
145 }
146 }
147 }
148}
149
150impl std::error::Error for RemoteMobError {}
151
152#[derive(Debug, Clone, PartialEq, Eq)]
154pub enum RemoteEndpoint {
155 Tcp(String),
157 Uds(String),
159}
160
161impl RemoteEndpoint {
162 pub fn scheme(&self) -> &'static str {
164 match self {
165 Self::Tcp(_) => "tcp",
166 Self::Uds(_) => "uds",
167 }
168 }
169
170 pub fn comms_address(&self) -> String {
173 match self {
174 Self::Tcp(addr) => format!("tcp://{addr}"),
175 Self::Uds(path) => format!("uds://{path}"),
176 }
177 }
178
179 pub fn raw(&self) -> &str {
181 match self {
182 Self::Tcp(s) | Self::Uds(s) => s.as_str(),
183 }
184 }
185}
186
187#[derive(Debug, Clone)]
194pub struct RemoteMobProxy {
195 mob_id: String,
196 endpoint: RemoteEndpoint,
197}
198
199impl RemoteMobProxy {
200 pub fn from_entry(entry: &ContactEntry) -> Result<Option<Self>, RemoteMobError> {
204 match &entry.transport {
205 MobTransport::Inproc => Ok(None),
206 MobTransport::Tcp(addr) => Ok(Some(Self {
207 mob_id: entry.mob_id.clone(),
208 endpoint: RemoteEndpoint::Tcp(addr.clone()),
209 })),
210 MobTransport::Uds(path) => Ok(Some(Self {
211 mob_id: entry.mob_id.clone(),
212 endpoint: RemoteEndpoint::Uds(path.clone()),
213 })),
214 }
215 }
216
217 pub fn mob_id(&self) -> &str {
219 &self.mob_id
220 }
221
222 pub fn endpoint(&self) -> &RemoteEndpoint {
224 &self.endpoint
225 }
226
227 pub async fn wire_remote(
234 &self,
235 remote_member: &str,
236 local_peer_spec_address: &str,
237 local_comms_name: &str,
238 local_peer_id: &str,
239 local_pubkey_b64: Option<String>,
240 ) -> Result<(), RemoteMobError> {
241 let request = super::cross_mob_control::ControlRequest::Wire {
242 remote_member: remote_member.to_string(),
243 local_peer_spec_address: local_peer_spec_address.to_string(),
244 local_comms_name: local_comms_name.to_string(),
245 local_peer_id: local_peer_id.to_string(),
246 local_pubkey_b64,
247 };
248 self.dispatch_no_payload(request, "wire").await
249 }
250
251 pub async fn unwire_remote(
253 &self,
254 remote_member: &str,
255 local_peer_spec_address: &str,
256 local_comms_name: &str,
257 local_peer_id: &str,
258 local_pubkey_b64: Option<String>,
259 ) -> Result<(), RemoteMobError> {
260 let request = super::cross_mob_control::ControlRequest::Unwire {
261 remote_member: remote_member.to_string(),
262 local_peer_spec_address: local_peer_spec_address.to_string(),
263 local_comms_name: local_comms_name.to_string(),
264 local_peer_id: local_peer_id.to_string(),
265 local_pubkey_b64,
266 };
267 self.dispatch_no_payload(request, "unwire").await
268 }
269
270 pub async fn inject_message(
276 &self,
277 remote_member: &str,
278 content_json: serde_json::Value,
279 ) -> Result<String, RemoteMobError> {
280 let request = super::cross_mob_control::ControlRequest::Inject {
281 remote_member: remote_member.to_string(),
282 content: content_json,
283 };
284 let response = super::cross_mob_control::RemoteControlClient::send(
285 &self.endpoint,
286 &request,
287 super::cross_mob_control::DEFAULT_CONTROL_TIMEOUT,
288 )
289 .await
290 .map_err(|err| self.attach_mob_id(err))?;
291 match response {
292 super::cross_mob_control::ControlResponse::Injected { session_id } => Ok(session_id),
293 super::cross_mob_control::ControlResponse::Err { code, message } => {
294 Err(RemoteMobError::Rejected {
295 mob_id: self.mob_id.clone(),
296 endpoint: self.endpoint.comms_address(),
297 code,
298 message,
299 })
300 }
301 other => Err(RemoteMobError::Decode {
302 endpoint: self.endpoint.comms_address(),
303 message: format!("expected Injected response, got {other:?}"),
304 }),
305 }
306 }
307
308 async fn dispatch_no_payload(
309 &self,
310 request: super::cross_mob_control::ControlRequest,
311 operation: &'static str,
312 ) -> Result<(), RemoteMobError> {
313 let response = super::cross_mob_control::RemoteControlClient::send(
314 &self.endpoint,
315 &request,
316 super::cross_mob_control::DEFAULT_CONTROL_TIMEOUT,
317 )
318 .await
319 .map_err(|err| self.attach_mob_id(err))?;
320 match response {
321 super::cross_mob_control::ControlResponse::Ok => Ok(()),
322 super::cross_mob_control::ControlResponse::Err { code, message } => {
323 Err(RemoteMobError::Rejected {
324 mob_id: self.mob_id.clone(),
325 endpoint: self.endpoint.comms_address(),
326 code,
327 message,
328 })
329 }
330 other => Err(RemoteMobError::Decode {
331 endpoint: self.endpoint.comms_address(),
332 message: format!("expected Ok for {operation}, got {other:?}"),
333 }),
334 }
335 }
336
337 pub async fn lookup_member(
341 &self,
342 remote_member: &str,
343 ) -> Result<(String, String), RemoteMobError> {
344 let request = super::cross_mob_control::ControlRequest::LookupMember {
345 remote_member: remote_member.to_string(),
346 };
347 let response = super::cross_mob_control::RemoteControlClient::send(
348 &self.endpoint,
349 &request,
350 super::cross_mob_control::DEFAULT_CONTROL_TIMEOUT,
351 )
352 .await
353 .map_err(|err| self.attach_mob_id(err))?;
354 match response {
355 super::cross_mob_control::ControlResponse::Member {
356 peer_id,
357 comms_name,
358 } => Ok((peer_id, comms_name)),
359 super::cross_mob_control::ControlResponse::Err { code, message } => {
360 Err(RemoteMobError::Rejected {
361 mob_id: self.mob_id.clone(),
362 endpoint: self.endpoint.comms_address(),
363 code,
364 message,
365 })
366 }
367 other => Err(RemoteMobError::Decode {
368 endpoint: self.endpoint.comms_address(),
369 message: format!("expected Member response, got {other:?}"),
370 }),
371 }
372 }
373
374 fn attach_mob_id(&self, err: RemoteMobError) -> RemoteMobError {
375 match err {
376 RemoteMobError::ControlChannelUnavailable {
377 mob_id,
378 endpoint,
379 operation,
380 } if mob_id.is_empty() || mob_id != self.mob_id => {
381 RemoteMobError::ControlChannelUnavailable {
382 mob_id: self.mob_id.clone(),
383 endpoint,
384 operation,
385 }
386 }
387 other => other,
388 }
389 }
390}
391
392#[cfg(test)]
393#[allow(clippy::unwrap_used, clippy::expect_used)]
394mod tests {
395 use super::*;
396
397 #[test]
398 fn from_entry_inproc_returns_none() {
399 let entry = ContactEntry {
400 mob_id: "demo".to_string(),
401 transport: MobTransport::Inproc,
402 pubkey: None,
403 };
404 let proxy = RemoteMobProxy::from_entry(&entry).expect("inproc is supported");
405 assert!(proxy.is_none());
406 }
407
408 #[test]
409 fn from_entry_tcp_round_trip() {
410 let entry = ContactEntry {
411 mob_id: "remote".to_string(),
412 transport: MobTransport::Tcp("127.0.0.1:9001".to_string()),
413 pubkey: None,
414 };
415 let proxy = RemoteMobProxy::from_entry(&entry)
416 .expect("tcp is supported")
417 .expect("tcp returns Some(proxy)");
418 assert_eq!(proxy.mob_id(), "remote");
419 assert_eq!(proxy.endpoint().scheme(), "tcp");
420 assert_eq!(proxy.endpoint().raw(), "127.0.0.1:9001");
421 assert_eq!(proxy.endpoint().comms_address(), "tcp://127.0.0.1:9001");
422 }
423
424 #[test]
425 fn from_entry_uds_round_trip() {
426 let entry = ContactEntry {
427 mob_id: "remote-uds".to_string(),
428 transport: MobTransport::Uds("/tmp/cross-mob.sock".to_string()),
429 pubkey: None,
430 };
431 let proxy = RemoteMobProxy::from_entry(&entry)
432 .expect("uds is supported")
433 .expect("uds returns Some(proxy)");
434 assert_eq!(proxy.endpoint().scheme(), "uds");
435 assert_eq!(proxy.endpoint().raw(), "/tmp/cross-mob.sock");
436 assert_eq!(
437 proxy.endpoint().comms_address(),
438 "uds:///tmp/cross-mob.sock"
439 );
440 }
441
442 #[tokio::test]
447 async fn wire_remote_returns_unavailable_when_no_listener() {
448 let entry = ContactEntry {
449 mob_id: "remote".to_string(),
450 transport: MobTransport::Tcp("127.0.0.1:1".to_string()),
451 pubkey: None,
452 };
453 let proxy = RemoteMobProxy::from_entry(&entry)
454 .expect("tcp ok")
455 .expect("some");
456 let err = proxy
457 .wire_remote(
458 "alice",
459 "tcp://127.0.0.1:9000",
460 "demo/role/alice",
461 "00000000-0000-4000-8000-000000000001",
462 None,
463 )
464 .await
465 .expect_err("no listener");
466 assert!(
467 matches!(err, RemoteMobError::ControlChannelUnavailable { ref mob_id, .. } if mob_id == "remote"),
468 "got {err:?}"
469 );
470 }
471}