1use std::ffi::c_char;
15use std::ptr;
16use std::sync::Arc;
17
18use pamoja_core::Transport;
19use pamoja_loopback::{LoopbackBroker, LoopbackTransport};
20use tokio::sync::Mutex;
21
22use crate::transport::{status, Kind, PamojaMessage, PamojaTransport};
23use crate::{read_bytes, read_str, runtime, set_last_error, PamojaStatus};
24
25pub struct PamojaLoopbackBroker {
30 inner: LoopbackBroker,
31}
32
33pub struct PamojaLoopbackTransport {
35 inner: Arc<Mutex<LoopbackTransport>>,
36}
37
38#[no_mangle]
44pub extern "C" fn pamoja_loopback_broker_new() -> *mut PamojaLoopbackBroker {
45 Box::into_raw(Box::new(PamojaLoopbackBroker {
46 inner: LoopbackBroker::new(),
47 }))
48}
49
50#[no_mangle]
62pub unsafe extern "C" fn pamoja_loopback_broker_free(broker: *mut PamojaLoopbackBroker) {
63 if !broker.is_null() {
64 drop(Box::from_raw(broker));
65 }
66}
67
68#[no_mangle]
84pub unsafe extern "C" fn pamoja_loopback_transport_new(
85 broker: *const PamojaLoopbackBroker,
86) -> *mut PamojaLoopbackTransport {
87 let Some(broker) = broker_handle(broker) else {
88 return ptr::null_mut();
89 };
90 Box::into_raw(Box::new(PamojaLoopbackTransport {
91 inner: Arc::new(Mutex::new(LoopbackTransport::new(broker.inner.clone()))),
92 }))
93}
94
95#[no_mangle]
109pub unsafe extern "C" fn pamoja_loopback_transport_connect(
110 transport: *mut PamojaLoopbackTransport,
111) -> PamojaStatus {
112 let Some(transport) = transport_handle(transport) else {
113 return PamojaStatus::InvalidArgument;
114 };
115 let inner = Arc::clone(&transport.inner);
116 status(runtime().block_on(async move { inner.lock().await.connect().await }))
117}
118
119#[no_mangle]
138pub unsafe extern "C" fn pamoja_loopback_transport_send(
139 transport: *mut PamojaLoopbackTransport,
140 topic: *const c_char,
141 payload: *const u8,
142 payload_len: usize,
143) -> PamojaStatus {
144 let Some(transport) = transport_handle(transport) else {
145 return PamojaStatus::InvalidArgument;
146 };
147 let Some(topic) = read_str(topic, "topic") else {
148 return PamojaStatus::InvalidArgument;
149 };
150 let payload = match read_bytes(payload, payload_len) {
151 Ok(payload) => payload,
152 Err(status) => return status,
153 };
154 let inner = Arc::clone(&transport.inner);
155 let topic = topic.to_owned();
156 status(runtime().block_on(async move { inner.lock().await.send(&topic, &payload).await }))
157}
158
159#[no_mangle]
175pub unsafe extern "C" fn pamoja_loopback_transport_subscribe(
176 transport: *mut PamojaLoopbackTransport,
177 topic: *const c_char,
178) -> PamojaStatus {
179 let Some(transport) = transport_handle(transport) else {
180 return PamojaStatus::InvalidArgument;
181 };
182 let Some(topic) = read_str(topic, "topic") else {
183 return PamojaStatus::InvalidArgument;
184 };
185 let inner = Arc::clone(&transport.inner);
186 let topic = topic.to_owned();
187 status(runtime().block_on(async move { inner.lock().await.subscribe(&topic).await }))
188}
189
190#[no_mangle]
206pub unsafe extern "C" fn pamoja_loopback_transport_recv(
207 transport: *mut PamojaLoopbackTransport,
208 out_message: *mut *mut PamojaMessage,
209) -> PamojaStatus {
210 let Some(transport) = transport_handle(transport) else {
211 return PamojaStatus::InvalidArgument;
212 };
213 if out_message.is_null() {
214 set_last_error("out_message must not be null".to_owned());
215 return PamojaStatus::InvalidArgument;
216 }
217 *out_message = ptr::null_mut();
218
219 let inner = Arc::clone(&transport.inner);
220 match runtime().block_on(async move { inner.lock().await.recv().await }) {
221 Ok(Some(message)) => {
222 *out_message = PamojaMessage::into_raw(message.topic, message.payload);
223 PamojaStatus::Ok
224 }
225 Ok(None) => PamojaStatus::Ok,
226 Err(error) => {
227 let code = PamojaStatus::from_error(&error);
228 set_last_error(error.to_string());
229 code
230 }
231 }
232}
233
234#[no_mangle]
249pub unsafe extern "C" fn pamoja_loopback_transport_is_connected(
250 transport: *mut PamojaLoopbackTransport,
251) -> bool {
252 let Some(transport) = transport_handle(transport) else {
253 return false;
254 };
255 let inner = Arc::clone(&transport.inner);
256 runtime().block_on(async move { inner.lock().await.is_connected() })
257}
258
259#[no_mangle]
270pub unsafe extern "C" fn pamoja_loopback_transport_disconnect(
271 transport: *mut PamojaLoopbackTransport,
272) {
273 let Some(transport) = transport_handle(transport) else {
274 return;
275 };
276 let inner = Arc::clone(&transport.inner);
277 runtime().block_on(async move { inner.lock().await.disconnect() });
278}
279
280#[no_mangle]
289pub unsafe extern "C" fn pamoja_loopback_transport_free(transport: *mut PamojaLoopbackTransport) {
290 if !transport.is_null() {
291 drop(Box::from_raw(transport));
292 }
293}
294
295#[no_mangle]
311pub unsafe extern "C" fn pamoja_transport_loopback(
312 broker: *const PamojaLoopbackBroker,
313) -> *mut PamojaTransport {
314 let Some(broker) = broker_handle(broker) else {
315 return ptr::null_mut();
316 };
317 PamojaTransport::into_raw(Kind::Loopback(LoopbackTransport::new(broker.inner.clone())))
318}
319
320unsafe fn broker_handle<'a>(
326 broker: *const PamojaLoopbackBroker,
327) -> Option<&'a PamojaLoopbackBroker> {
328 if broker.is_null() {
329 set_last_error("broker must not be null".to_owned());
330 return None;
331 }
332 Some(&*broker)
333}
334
335unsafe fn transport_handle<'a>(
342 transport: *mut PamojaLoopbackTransport,
343) -> Option<&'a PamojaLoopbackTransport> {
344 if transport.is_null() {
345 set_last_error("transport must not be null".to_owned());
346 return None;
347 }
348 Some(&*transport)
349}
350
351#[cfg(test)]
352mod tests {
353 use super::*;
354 use crate::transport::{
355 pamoja_message_free, pamoja_message_payload, pamoja_message_payload_len,
356 };
357
358 unsafe fn take(message: *mut PamojaMessage) -> Vec<u8> {
360 assert!(!message.is_null());
361 let bytes = std::slice::from_raw_parts(
362 pamoja_message_payload(message),
363 pamoja_message_payload_len(message),
364 )
365 .to_vec();
366 pamoja_message_free(message);
367 bytes
368 }
369
370 #[test]
371 fn a_published_message_reaches_a_subscriber() {
372 unsafe {
373 let broker = pamoja_loopback_broker_new();
374 let publisher = pamoja_loopback_transport_new(broker);
375 let subscriber = pamoja_loopback_transport_new(broker);
376
377 assert_eq!(
378 pamoja_loopback_transport_connect(publisher),
379 PamojaStatus::Ok
380 );
381 assert_eq!(
382 pamoja_loopback_transport_connect(subscriber),
383 PamojaStatus::Ok
384 );
385 assert!(pamoja_loopback_transport_is_connected(publisher));
386
387 let topic = std::ffi::CString::new("sensors/1").expect("static");
388 assert_eq!(
389 pamoja_loopback_transport_subscribe(subscriber, topic.as_ptr()),
390 PamojaStatus::Ok
391 );
392 assert_eq!(
393 pamoja_loopback_transport_send(publisher, topic.as_ptr(), b"21.5".as_ptr(), 4),
394 PamojaStatus::Ok
395 );
396
397 let mut message = ptr::null_mut();
398 assert_eq!(
399 pamoja_loopback_transport_recv(subscriber, &mut message),
400 PamojaStatus::Ok
401 );
402 assert_eq!(take(message), b"21.5");
403
404 pamoja_loopback_transport_free(subscriber);
405 pamoja_loopback_transport_free(publisher);
406 pamoja_loopback_broker_free(broker);
407 }
408 }
409
410 #[test]
411 fn a_disconnected_link_refuses_to_send() {
412 unsafe {
413 let broker = pamoja_loopback_broker_new();
414 let transport = pamoja_loopback_transport_new(broker);
415 pamoja_loopback_transport_connect(transport);
416 pamoja_loopback_transport_disconnect(transport);
417 assert!(!pamoja_loopback_transport_is_connected(transport));
418
419 let topic = std::ffi::CString::new("sensors/1").expect("static");
420 assert_ne!(
421 pamoja_loopback_transport_send(transport, topic.as_ptr(), b"x".as_ptr(), 1),
422 PamojaStatus::Ok
423 );
424
425 pamoja_loopback_transport_free(transport);
426 pamoja_loopback_broker_free(broker);
427 }
428 }
429
430 #[test]
431 fn a_null_handle_is_refused_rather_than_dereferenced() {
432 unsafe {
433 assert!(pamoja_loopback_transport_new(ptr::null()).is_null());
434 assert!(!pamoja_loopback_transport_is_connected(ptr::null_mut()));
435 assert_eq!(
436 pamoja_loopback_transport_connect(ptr::null_mut()),
437 PamojaStatus::InvalidArgument
438 );
439 pamoja_loopback_transport_free(ptr::null_mut());
440 pamoja_loopback_broker_free(ptr::null_mut());
441 }
442 }
443}