Skip to main content

pamoja_ffi/
ladder.rs

1//! The C ABI for the cost-aware transport ladder.
2//!
3//! These functions wrap [`pamoja_ladder`] for callers that reach the SDK through
4//! the flat C boundary. A ladder is the answer to a node that has more than one
5//! way to reach the network and no single one that always works: rungs are tried
6//! in the order they were added, cheapest first, and a message no rung accepts
7//! goes into a buffer rather than being lost.
8//!
9//! The Rust ladder is generic over its buffer, which cannot cross a C ABI, so
10//! this one is built over the store handle from [`crate::sync`]. That handle
11//! already covers both an in-memory and a file-backed queue, so nothing is given
12//! up: a caller still chooses whether the buffer survives a restart.
13
14use std::ffi::c_char;
15use std::ptr;
16
17use pamoja_ladder::{Delivery, TransportLadder};
18
19use crate::sync::{take_store, PamojaStore, StoreKind};
20use crate::transport::{take_transport, PamojaTransport};
21use crate::{read_bytes, read_str, runtime, set_last_error, PamojaStatus};
22
23/// What became of a message handed to a ladder.
24#[repr(C)]
25#[derive(Clone, Copy, Debug, PartialEq, Eq)]
26pub enum PamojaDelivery {
27    /// A rung took the message and it is on its way.
28    Sent = 0,
29    /// No rung would take it, so it is in the buffer awaiting a flush.
30    Buffered = 1,
31}
32
33/// An opaque handle to a ladder and the buffer behind it.
34///
35/// The Rust builder takes the ladder by value to add a rung, which would move
36/// the handle out from under a caller holding a pointer to it. Holding the
37/// ladder in an option lets it be taken and put back so the handle address stays
38/// good for the life of the ladder.
39pub struct PamojaLadder {
40    inner: Option<TransportLadder<StoreKind>>,
41}
42
43/// Creates a ladder with no rungs, buffering into a store.
44///
45/// # Arguments
46///
47/// * `store` - the buffer to hold messages no rung would take, consumed by this
48///   call.
49///
50/// # Returns
51///
52/// A handle the caller must release with [`pamoja_ladder_free`], or null if
53/// `store` is null.
54///
55/// # Safety
56///
57/// `store` must be a live handle from [`crate::sync`] that has not been freed or
58/// consumed. After this call it must not be used again, whatever the result.
59#[no_mangle]
60pub unsafe extern "C" fn pamoja_ladder_new(store: *mut PamojaStore) -> *mut PamojaLadder {
61    let Some(store) = take_store(store) else {
62        return ptr::null_mut();
63    };
64    Box::into_raw(Box::new(PamojaLadder {
65        inner: Some(TransportLadder::new(store)),
66    }))
67}
68
69/// Adds a rung, which is tried after the rungs already added.
70///
71/// Add the cheapest, most-preferred link first and the costliest fallback last,
72/// because a send takes the first rung that accepts it.
73///
74/// # Arguments
75///
76/// * `ladder` - the ladder to add to.
77/// * `transport` - the transport to add, consumed by this call.
78///
79/// # Returns
80///
81/// [`PamojaStatus::Ok`] once the rung is added.
82///
83/// # Safety
84///
85/// `ladder` must be a live handle, and `transport` a live transport handle that
86/// has not been freed or consumed. After this call the transport must not be
87/// used again, whatever the result.
88#[no_mangle]
89pub unsafe extern "C" fn pamoja_ladder_rung(
90    ladder: *mut PamojaLadder,
91    transport: *mut PamojaTransport,
92) -> PamojaStatus {
93    let Some(handle) = ladder_handle(ladder) else {
94        // The transport was promised to this call, so it is released rather than
95        // left for a caller who has already been told not to touch it again.
96        drop(take_transport(transport));
97        return PamojaStatus::InvalidArgument;
98    };
99    let Some(transport) = take_transport(transport) else {
100        return PamojaStatus::InvalidArgument;
101    };
102    let Some(inner) = handle.inner.take() else {
103        set_last_error("this ladder is no longer usable".to_owned());
104        return PamojaStatus::InvalidArgument;
105    };
106    handle.inner = Some(inner.rung(transport));
107    PamojaStatus::Ok
108}
109
110/// Connects every rung, so a send can be tried against each in turn.
111///
112/// A rung that will not connect is left in the ladder: it may come back, and a
113/// send simply falls through it until it does.
114///
115/// # Arguments
116///
117/// * `ladder` - the ladder.
118///
119/// # Returns
120///
121/// [`PamojaStatus::Ok`] once the rungs have been tried.
122///
123/// # Safety
124///
125/// `ladder` must be a live handle from [`pamoja_ladder_new`].
126#[no_mangle]
127pub unsafe extern "C" fn pamoja_ladder_connect(ladder: *mut PamojaLadder) -> PamojaStatus {
128    let Some(inner) = ladder_inner(ladder) else {
129        return PamojaStatus::InvalidArgument;
130    };
131    match runtime().block_on(inner.connect()) {
132        Ok(()) => PamojaStatus::Ok,
133        Err(error) => fail(error),
134    }
135}
136
137/// Sends a payload, falling through the rungs and buffering if none take it.
138///
139/// # Arguments
140///
141/// * `ladder` - the ladder.
142/// * `topic` - the destination topic, as null-terminated UTF-8.
143/// * `payload` - the bytes to send.
144/// * `payload_len` - the length of `payload`.
145/// * `out_delivery` - receives whether the message went out or was buffered.
146///
147/// # Returns
148///
149/// [`PamojaStatus::Ok`] on success, with `out_delivery` saying which happened.
150/// Buffering is a success, not a failure: it is what the ladder exists to do.
151///
152/// # Safety
153///
154/// `ladder` must be a live handle, `topic` a valid null-terminated UTF-8 string,
155/// `payload` must point to at least `payload_len` readable bytes or be null when
156/// that length is 0, and `out_delivery` must be writable or null.
157#[no_mangle]
158pub unsafe extern "C" fn pamoja_ladder_send(
159    ladder: *mut PamojaLadder,
160    topic: *const c_char,
161    payload: *const u8,
162    payload_len: usize,
163    out_delivery: *mut PamojaDelivery,
164) -> PamojaStatus {
165    let Some(topic) = read_str(topic, "topic") else {
166        return PamojaStatus::InvalidArgument;
167    };
168    let payload = match read_bytes(payload, payload_len) {
169        Ok(payload) => payload,
170        Err(status) => return status,
171    };
172    let Some(inner) = ladder_inner(ladder) else {
173        return PamojaStatus::InvalidArgument;
174    };
175    match runtime().block_on(inner.send(topic, &payload)) {
176        Ok(delivery) => {
177            if !out_delivery.is_null() {
178                *out_delivery = match delivery {
179                    Delivery::Sent => PamojaDelivery::Sent,
180                    Delivery::Buffered => PamojaDelivery::Buffered,
181                };
182            }
183            PamojaStatus::Ok
184        }
185        Err(error) => fail(error),
186    }
187}
188
189/// Replays the buffer over the rungs, oldest message first.
190///
191/// # Arguments
192///
193/// * `ladder` - the ladder.
194/// * `out_sent` - receives how many buffered messages went out, or may be null.
195///
196/// # Returns
197///
198/// [`PamojaStatus::Ok`] on success.
199///
200/// # Safety
201///
202/// `ladder` must be a live handle and `out_sent` writable or null.
203#[no_mangle]
204pub unsafe extern "C" fn pamoja_ladder_flush(
205    ladder: *mut PamojaLadder,
206    out_sent: *mut usize,
207) -> PamojaStatus {
208    let Some(inner) = ladder_inner(ladder) else {
209        return PamojaStatus::InvalidArgument;
210    };
211    match runtime().block_on(inner.flush()) {
212        Ok(sent) => {
213            if !out_sent.is_null() {
214                *out_sent = sent;
215            }
216            PamojaStatus::Ok
217        }
218        Err(error) => fail(error),
219    }
220}
221
222/// Reports how many messages are waiting in the buffer.
223///
224/// # Arguments
225///
226/// * `ladder` - the ladder.
227/// * `out_count` - receives the count.
228///
229/// # Returns
230///
231/// [`PamojaStatus::Ok`] on success.
232///
233/// # Safety
234///
235/// `ladder` must be a live handle and `out_count` must be writable.
236#[no_mangle]
237pub unsafe extern "C" fn pamoja_ladder_buffered(
238    ladder: *mut PamojaLadder,
239    out_count: *mut usize,
240) -> PamojaStatus {
241    let Some(inner) = ladder_inner(ladder) else {
242        return PamojaStatus::InvalidArgument;
243    };
244    if out_count.is_null() {
245        set_last_error("out_count must not be null".to_owned());
246        return PamojaStatus::InvalidArgument;
247    }
248    match runtime().block_on(inner.buffered()) {
249        Ok(count) => {
250            *out_count = count;
251            PamojaStatus::Ok
252        }
253        Err(error) => fail(error),
254    }
255}
256
257/// Releases a ladder handle, and the rungs and buffer it owns.
258///
259/// Passing null is a no-op.
260///
261/// # Safety
262///
263/// `ladder` must be a handle from [`pamoja_ladder_new`] that has not already
264/// been freed, or null. After this call it must not be used again.
265#[no_mangle]
266pub unsafe extern "C" fn pamoja_ladder_free(ladder: *mut PamojaLadder) {
267    if !ladder.is_null() {
268        drop(Box::from_raw(ladder));
269    }
270}
271
272/// Borrows a ladder handle, rejecting a null pointer.
273///
274/// # Safety
275///
276/// `ladder` must be a live handle from [`pamoja_ladder_new`], or null.
277unsafe fn ladder_handle<'a>(ladder: *mut PamojaLadder) -> Option<&'a mut PamojaLadder> {
278    if ladder.is_null() {
279        set_last_error("ladder must not be null".to_owned());
280        return None;
281    }
282    Some(&mut *ladder)
283}
284
285/// Borrows the ladder inside a handle, rejecting a null or spent one.
286///
287/// # Safety
288///
289/// `ladder` must be a live handle from [`pamoja_ladder_new`], or null.
290unsafe fn ladder_inner<'a>(
291    ladder: *mut PamojaLadder,
292) -> Option<&'a mut TransportLadder<StoreKind>> {
293    let handle = ladder_handle(ladder)?;
294    match handle.inner.as_mut() {
295        Some(inner) => Some(inner),
296        None => {
297            set_last_error("this ladder is no longer usable".to_owned());
298            None
299        }
300    }
301}
302
303/// Records an error and maps it onto a status.
304fn fail(error: pamoja_core::Error) -> PamojaStatus {
305    let status = PamojaStatus::from_error(&error);
306    set_last_error(error.to_string());
307    status
308}
309
310#[cfg(test)]
311mod tests {
312    use super::*;
313    use crate::loopback::{
314        pamoja_loopback_broker_free, pamoja_loopback_broker_new, pamoja_loopback_transport_connect,
315        pamoja_loopback_transport_free, pamoja_loopback_transport_new,
316        pamoja_loopback_transport_recv, pamoja_loopback_transport_subscribe,
317        pamoja_transport_loopback,
318    };
319    use crate::sync::pamoja_store_memory;
320    use crate::transport::{
321        pamoja_message_free, pamoja_message_payload, pamoja_message_payload_len,
322        pamoja_transport_faulty,
323    };
324
325    #[test]
326    fn a_message_no_rung_takes_is_buffered_rather_than_lost() {
327        unsafe {
328            let ladder = pamoja_ladder_new(pamoja_store_memory(0));
329            let topic = std::ffi::CString::new("sensors/1").expect("static");
330
331            let mut delivery = PamojaDelivery::Sent;
332            assert_eq!(
333                pamoja_ladder_send(ladder, topic.as_ptr(), b"21.5".as_ptr(), 4, &mut delivery),
334                PamojaStatus::Ok,
335                "buffering is a success, not a failure"
336            );
337            assert_eq!(delivery, PamojaDelivery::Buffered);
338
339            let mut waiting = 0;
340            assert_eq!(
341                pamoja_ladder_buffered(ladder, &mut waiting),
342                PamojaStatus::Ok
343            );
344            assert_eq!(waiting, 1);
345
346            pamoja_ladder_free(ladder);
347        }
348    }
349
350    #[test]
351    fn a_rung_that_refuses_falls_through_to_the_next() {
352        unsafe {
353            let broker = pamoja_loopback_broker_new();
354            let listener = pamoja_loopback_transport_new(broker);
355            pamoja_loopback_transport_connect(listener);
356            let topic = std::ffi::CString::new("sensors/1").expect("static");
357            pamoja_loopback_transport_subscribe(listener, topic.as_ptr());
358
359            // The first rung fails its next send; the second is the same broker.
360            let failing = pamoja_transport_faulty(pamoja_transport_loopback(broker), 1);
361            let working = pamoja_transport_loopback(broker);
362
363            let ladder = pamoja_ladder_new(pamoja_store_memory(0));
364            assert_eq!(pamoja_ladder_rung(ladder, failing), PamojaStatus::Ok);
365            assert_eq!(pamoja_ladder_rung(ladder, working), PamojaStatus::Ok);
366            assert_eq!(pamoja_ladder_connect(ladder), PamojaStatus::Ok);
367
368            let mut delivery = PamojaDelivery::Buffered;
369            assert_eq!(
370                pamoja_ladder_send(ladder, topic.as_ptr(), b"21.5".as_ptr(), 4, &mut delivery),
371                PamojaStatus::Ok
372            );
373            assert_eq!(
374                delivery,
375                PamojaDelivery::Sent,
376                "the second rung carried what the first refused"
377            );
378
379            let mut message = ptr::null_mut();
380            pamoja_loopback_transport_recv(listener, &mut message);
381            assert!(!message.is_null());
382            let payload = std::slice::from_raw_parts(
383                pamoja_message_payload(message),
384                pamoja_message_payload_len(message),
385            )
386            .to_vec();
387            assert_eq!(payload, b"21.5");
388            pamoja_message_free(message);
389
390            pamoja_ladder_free(ladder);
391            pamoja_loopback_transport_free(listener);
392            pamoja_loopback_broker_free(broker);
393        }
394    }
395
396    #[test]
397    fn a_flush_replays_what_was_buffered() {
398        unsafe {
399            let broker = pamoja_loopback_broker_new();
400            let listener = pamoja_loopback_transport_new(broker);
401            pamoja_loopback_transport_connect(listener);
402            let topic = std::ffi::CString::new("sensors/1").expect("static");
403            pamoja_loopback_transport_subscribe(listener, topic.as_ptr());
404
405            let ladder = pamoja_ladder_new(pamoja_store_memory(0));
406
407            // Nothing to send over yet, so it buffers.
408            pamoja_ladder_send(ladder, topic.as_ptr(), b"one".as_ptr(), 3, ptr::null_mut());
409            pamoja_ladder_send(ladder, topic.as_ptr(), b"two".as_ptr(), 3, ptr::null_mut());
410
411            // The link comes back.
412            assert_eq!(
413                pamoja_ladder_rung(ladder, pamoja_transport_loopback(broker)),
414                PamojaStatus::Ok
415            );
416            assert_eq!(pamoja_ladder_connect(ladder), PamojaStatus::Ok);
417
418            let mut sent = 0;
419            assert_eq!(pamoja_ladder_flush(ladder, &mut sent), PamojaStatus::Ok);
420            assert_eq!(sent, 2, "both buffered messages went out");
421
422            let mut waiting = 1;
423            pamoja_ladder_buffered(ladder, &mut waiting);
424            assert_eq!(waiting, 0, "and the buffer is empty");
425
426            pamoja_ladder_free(ladder);
427            pamoja_loopback_transport_free(listener);
428            pamoja_loopback_broker_free(broker);
429        }
430    }
431
432    #[test]
433    fn a_null_handle_is_refused_rather_than_dereferenced() {
434        unsafe {
435            assert!(pamoja_ladder_new(ptr::null_mut()).is_null());
436            assert_eq!(
437                pamoja_ladder_connect(ptr::null_mut()),
438                PamojaStatus::InvalidArgument
439            );
440            assert_eq!(
441                pamoja_ladder_buffered(ptr::null_mut(), ptr::null_mut()),
442                PamojaStatus::InvalidArgument
443            );
444            pamoja_ladder_free(ptr::null_mut());
445        }
446    }
447}