1use 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#[repr(C)]
25#[derive(Clone, Copy, Debug, PartialEq, Eq)]
26pub enum PamojaDelivery {
27 Sent = 0,
29 Buffered = 1,
31}
32
33pub struct PamojaLadder {
40 inner: Option<TransportLadder<StoreKind>>,
41}
42
43#[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#[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 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#[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#[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#[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#[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#[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
272unsafe 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
285unsafe 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
303fn 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 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 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 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}