Skip to main content

pamoja_ffi/
loopback.rs

1//! The C ABI for the in-process loopback broker.
2//!
3//! These functions wrap [`pamoja_loopback`] so a caller can exercise the
4//! publish-and-subscribe path with no broker, no network, and no hardware. That
5//! matters most for the bindings: someone writing against the SDK from Python or
6//! C# can drive a whole message flow in a unit test rather than standing up
7//! infrastructure to find out whether their topics line up.
8//!
9//! A broker is shared by cloning, so every transport built from one sees the
10//! same traffic. Pair it with
11//! [`pamoja_transport_faulty`](crate::transport::pamoja_transport_faulty) to
12//! check what a caller does when a link starts refusing sends.
13
14use 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
25/// An opaque handle to an in-process broker.
26///
27/// Every transport built from one broker shares its traffic, so a message one
28/// publishes reaches the others that subscribed to the topic.
29pub struct PamojaLoopbackBroker {
30    inner: LoopbackBroker,
31}
32
33/// An opaque handle to one in-process link to a broker.
34pub struct PamojaLoopbackTransport {
35    inner: Arc<Mutex<LoopbackTransport>>,
36}
37
38/// Creates an in-process broker with no traffic.
39///
40/// # Returns
41///
42/// A handle the caller must release with [`pamoja_loopback_broker_free`].
43#[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/// Releases a broker handle.
51///
52/// Transports already built from the broker keep working, because each holds
53/// its own share of it.
54///
55/// Passing null is a no-op.
56///
57/// # Safety
58///
59/// `broker` must be a handle from [`pamoja_loopback_broker_new`] that has not
60/// already been freed, or null. After this call it must not be used again.
61#[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/// Creates a link to a broker.
69///
70/// # Arguments
71///
72/// * `broker` - the broker to join. It is shared, not consumed, so the same
73///   broker can back as many links as the caller needs.
74///
75/// # Returns
76///
77/// A handle the caller must release with
78/// [`pamoja_loopback_transport_free`], or null if `broker` is null.
79///
80/// # Safety
81///
82/// `broker` must be a live handle from [`pamoja_loopback_broker_new`], or null.
83#[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/// Marks a link connected so it will carry traffic.
96///
97/// # Arguments
98///
99/// * `transport` - the link.
100///
101/// # Returns
102///
103/// [`PamojaStatus::Ok`] once connected.
104///
105/// # Safety
106///
107/// `transport` must be a live handle from [`pamoja_loopback_transport_new`].
108#[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/// Publishes a payload to a topic on the broker.
120///
121/// # Arguments
122///
123/// * `transport` - the link.
124/// * `topic` - the destination topic, as null-terminated UTF-8.
125/// * `payload` - the bytes to publish.
126/// * `payload_len` - the length of `payload`.
127///
128/// # Returns
129///
130/// [`PamojaStatus::Ok`] once every subscriber has been handed the message.
131///
132/// # Safety
133///
134/// `transport` must be a live handle, `topic` a valid null-terminated UTF-8
135/// string, and `payload` must point to at least `payload_len` readable bytes or
136/// be null when that length is 0.
137#[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/// Subscribes a link to a topic.
160///
161/// # Arguments
162///
163/// * `transport` - the link.
164/// * `topic` - the topic to subscribe to, as null-terminated UTF-8.
165///
166/// # Returns
167///
168/// [`PamojaStatus::Ok`] once subscribed.
169///
170/// # Safety
171///
172/// `transport` must be a live handle and `topic` a valid null-terminated UTF-8
173/// string.
174#[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/// Waits for the next message on a subscribed topic.
191///
192/// # Arguments
193///
194/// * `transport` - the link.
195/// * `out_message` - receives a message handle, or null when the link is closed.
196///
197/// # Returns
198///
199/// [`PamojaStatus::Ok`] on success. A null `out_message` with an `Ok` status
200/// means the link closed rather than that anything failed.
201///
202/// # Safety
203///
204/// `transport` must be a live handle and `out_message` must be writable.
205#[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/// Reports whether a link is connected.
235///
236/// # Arguments
237///
238/// * `transport` - the link.
239///
240/// # Returns
241///
242/// `true` when connected, or `false` if `transport` is null.
243///
244/// # Safety
245///
246/// `transport` must be a live handle from [`pamoja_loopback_transport_new`], or
247/// null.
248#[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/// Marks a link disconnected, so sends over it fail.
260///
261/// # Arguments
262///
263/// * `transport` - the link.
264///
265/// # Safety
266///
267/// `transport` must be a live handle from [`pamoja_loopback_transport_new`], or
268/// null.
269#[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/// Releases a link handle.
281///
282/// Passing null is a no-op.
283///
284/// # Safety
285///
286/// `transport` must be a handle from [`pamoja_loopback_transport_new`] that has
287/// not already been freed, or null. After this call it must not be used again.
288#[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/// Creates a loopback transport for composing into a ladder or a wrapper.
296///
297/// # Arguments
298///
299/// * `broker` - the broker to join, shared rather than consumed.
300///
301/// # Returns
302///
303/// A handle the caller must release with
304/// [`pamoja_transport_free`](crate::transport::pamoja_transport_free) or hand to
305/// a call that consumes it, or null if `broker` is null.
306///
307/// # Safety
308///
309/// `broker` must be a live handle from [`pamoja_loopback_broker_new`], or null.
310#[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
320/// Borrows a broker handle, rejecting a null pointer.
321///
322/// # Safety
323///
324/// `broker` must be a live handle from [`pamoja_loopback_broker_new`], or null.
325unsafe 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
335/// Borrows a link handle, rejecting a null pointer.
336///
337/// # Safety
338///
339/// `transport` must be a live handle from [`pamoja_loopback_transport_new`], or
340/// null.
341unsafe 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    /// Reads a message handle out and releases it.
359    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}