use std::ffi::{CStr, CString, c_char, c_void};
use std::marker::PhantomData;
use std::panic::{AssertUnwindSafe, catch_unwind};
use std::ptr::NonNull;
use std::slice;
use std::sync::Mutex;
use std::time::{Duration, Instant};
use crate::error::{Error, Result, last_ffi_error};
const POSE_MIN_INTERVAL: Duration = Duration::from_millis(200);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Protocol {
Udp,
Quic,
Tcp,
}
impl Default for Protocol {
fn default() -> Self {
Protocol::Quic
}
}
impl Protocol {
fn as_raw(self) -> adamo_sys::adamo_protocol_t {
match self {
Protocol::Udp => adamo_sys::ADAMO_PROTOCOL_UDP,
Protocol::Quic => adamo_sys::ADAMO_PROTOCOL_QUIC,
Protocol::Tcp => adamo_sys::ADAMO_PROTOCOL_TCP,
}
}
}
pub struct Session {
raw: NonNull<adamo_sys::adamo_session_t>,
pose_last_pub: Mutex<Option<Instant>>,
}
unsafe impl Send for Session {}
unsafe impl Sync for Session {}
impl Session {
pub(crate) fn raw_ptr(&self) -> *const adamo_sys::adamo_session_t {
self.raw.as_ptr()
}
pub fn open(api_key: &str, protocol: Protocol) -> Result<Self> {
let key = CString::new(api_key)?;
let raw = unsafe { adamo_sys::adamo_open(key.as_ptr(), protocol.as_raw()) };
match NonNull::new(raw) {
Some(raw) => Ok(Session {
raw,
pose_last_pub: Mutex::new(None),
}),
None => Err(last_ffi_error()),
}
}
pub fn open_mtls(api_key: &str, protocol: Protocol) -> Result<Self> {
let key = CString::new(api_key)?;
let raw = unsafe { adamo_sys::adamo_open_mtls(key.as_ptr(), protocol.as_raw()) };
match NonNull::new(raw) {
Some(raw) => Ok(Session {
raw,
pose_last_pub: Mutex::new(None),
}),
None => Err(last_ffi_error()),
}
}
pub fn open_default(api_key: &str) -> Result<Self> {
let key = CString::new(api_key)?;
let raw = unsafe { adamo_sys::adamo_open_default(key.as_ptr()) };
match NonNull::new(raw) {
Some(raw) => Ok(Session {
raw,
pose_last_pub: Mutex::new(None),
}),
None => Err(last_ffi_error()),
}
}
pub fn org(&self) -> Result<&str> {
let ptr = unsafe { adamo_sys::adamo_session_org(self.raw.as_ptr()) };
if ptr.is_null() {
return Err(last_ffi_error());
}
let len = unsafe { adamo_sys::adamo_session_org_len(self.raw.as_ptr()) };
let bytes = unsafe { slice::from_raw_parts(ptr.cast::<u8>(), len) };
std::str::from_utf8(bytes).map_err(|_| Error::InvalidUtf8)
}
pub fn relay_rtt(&self) -> Result<Option<Duration>> {
let mut out_us: u64 = 0;
let rc = unsafe { adamo_sys::adamo_relay_rtt_us(self.raw.as_ptr(), &mut out_us) };
if rc == 0 {
return Ok(Some(Duration::from_micros(out_us)));
}
match last_ffi_error() {
Error::Ffi(m) if m.contains("not available yet") => Ok(None),
e => Err(e),
}
}
pub fn put(&self, key: &str, payload: &[u8], opts: PublishOptions) -> Result<()> {
let key = CString::new(key)?;
let rc = unsafe {
adamo_sys::adamo_put(
self.raw.as_ptr(),
key.as_ptr(),
payload.as_ptr(),
payload.len(),
opts.priority,
opts.express as i32,
)
};
if rc == 0 { Ok(()) } else { Err(last_ffi_error()) }
}
pub fn state_store(&self, robot: &str) -> Result<StateStore<'_>> {
let robot = CString::new(robot)?;
let raw = unsafe { adamo_sys::adamo_state_declare(self.raw.as_ptr(), robot.as_ptr()) };
match NonNull::new(raw) {
Some(raw) => Ok(StateStore {
raw,
_session: PhantomData,
}),
None => Err(last_ffi_error()),
}
}
pub fn task_runner(&self, robot: &str) -> Result<TaskRunner<'_>> {
let robot = CString::new(robot)?;
let raw = unsafe { adamo_sys::adamo_task_runner_declare(self.raw.as_ptr(), robot.as_ptr()) };
match NonNull::new(raw) {
Some(raw) => Ok(TaskRunner {
raw,
_session: PhantomData,
}),
None => Err(last_ffi_error()),
}
}
pub fn publisher(&self, key: &str, opts: PublisherOptions) -> Result<Publisher<'_>> {
let key = CString::new(key)?;
let raw = unsafe {
adamo_sys::adamo_publisher(
self.raw.as_ptr(),
key.as_ptr(),
opts.priority,
opts.express as i32,
opts.reliable as i32,
)
};
match NonNull::new(raw) {
Some(raw) => Ok(Publisher {
raw,
_session: PhantomData,
}),
None => Err(last_ffi_error()),
}
}
pub fn log(&self, name: &str, message: &str, level: &str) -> Result<()> {
let key = format!("{name}/logs");
let truncated: String;
let message_ref = if message.len() > 10_000 {
let cutoff = message
.char_indices()
.take_while(|(i, _)| *i < 10_000)
.last()
.map_or(0, |(i, c)| i + c.len_utf8());
truncated = format!("{}... [truncated]", &message[..cutoff]);
truncated.as_str()
} else {
message
};
let payload = format!(
"{{\"ts_us\":{},\"level\":\"{}\",\"message\":{}}}",
crate::fabric_now_us(),
json_escape(level),
json_quote(message_ref),
);
self.put(
&key,
payload.as_bytes(),
PublishOptions {
priority: 230,
express: true,
},
)
}
pub fn set_pose(
&self,
name: &str,
x: f64,
y: f64,
z: f64,
frame: &str,
heading: Option<f64>,
) -> Result<()> {
if !(x.is_finite() && y.is_finite() && z.is_finite())
|| heading.is_some_and(|h| !h.is_finite())
{
return Err(Error::Invalid(
"set_pose: coordinates must be finite".into(),
));
}
{
let mut last = self.pose_last_pub.lock().unwrap();
let now = Instant::now();
if last.is_some_and(|t| now.duration_since(t) < POSE_MIN_INTERVAL) {
return Ok(());
}
*last = Some(now);
}
let heading_field = heading
.map(|h| format!(",\"heading\":{h}"))
.unwrap_or_default();
let payload = format!(
"{{\"ts_us\":{},\"frame\":{},\"x\":{x},\"y\":{y},\"z\":{z}{heading_field}}}",
crate::fabric_now_us(),
json_quote(frame),
);
self.put(
&format!("{name}/pose"),
payload.as_bytes(),
PublishOptions {
priority: Priority::DATA,
express: false,
},
)
}
pub fn subscribe(&self, key: &str) -> Result<Subscriber<'_>> {
let key = CString::new(key)?;
let raw = unsafe { adamo_sys::adamo_subscribe(self.raw.as_ptr(), key.as_ptr()) };
match NonNull::new(raw) {
Some(raw) => Ok(Subscriber {
raw,
_session: PhantomData,
}),
None => Err(last_ffi_error()),
}
}
pub fn subscribe_with<F>(&self, key: &str, callback: F) -> Result<CallbackSubscriber<'_>>
where
F: Fn(Sample) + Send + Sync + 'static,
{
let key = CString::new(key)?;
let mut state = Box::new(CallbackState {
callback: Box::new(callback),
});
let raw = unsafe {
adamo_sys::adamo_subscribe_cb(
self.raw.as_ptr(),
key.as_ptr(),
Some(callback_trampoline),
state.as_mut() as *mut CallbackState as *mut c_void,
)
};
match NonNull::new(raw) {
Some(raw) => Ok(CallbackSubscriber {
raw,
_state: state,
_session: PhantomData,
}),
None => Err(last_ffi_error()),
}
}
pub fn get(&self, key: &str, timeout: Duration) -> Result<Vec<Sample>> {
let key = CString::new(key)?;
let mut count = 0usize;
let raw = unsafe {
adamo_sys::adamo_get(
self.raw.as_ptr(),
key.as_ptr(),
timeout.as_millis().min(u64::MAX as u128) as u64,
&mut count,
)
};
if raw.is_null() {
return match last_ffi_error() {
Error::Ffi(m) if m == "(no error message)" && count == 0 => Ok(Vec::new()),
e => Err(e),
};
}
let raw_samples = unsafe { slice::from_raw_parts(raw, count) };
let mut samples = Vec::with_capacity(count);
for sample in raw_samples {
if let Some(sample) = Sample::from_borrowed(*sample) {
samples.push(sample);
}
}
unsafe { adamo_sys::adamo_get_replies_free(raw, count) };
Ok(samples)
}
pub fn alive(&self, token_key: &str) -> Result<LivelinessToken<'_>> {
let token_key = CString::new(token_key)?;
let raw = unsafe {
adamo_sys::adamo_liveliness_declare(self.raw.as_ptr(), token_key.as_ptr())
};
match NonNull::new(raw) {
Some(raw) => Ok(LivelinessToken {
raw,
_session: PhantomData,
}),
None => Err(last_ffi_error()),
}
}
pub fn live_tokens(&self, pattern: &str) -> Result<Vec<String>> {
let pattern = CString::new(pattern)?;
let mut count = 0usize;
let raw = unsafe {
adamo_sys::adamo_liveliness_get(self.raw.as_ptr(), pattern.as_ptr(), &mut count)
};
if raw.is_null() {
return match last_ffi_error() {
Error::Ffi(m) if m == "(no error message)" && count == 0 => Ok(Vec::new()),
e => Err(e),
};
}
let raw_tokens = unsafe { slice::from_raw_parts(raw, count) };
let mut tokens = Vec::with_capacity(count);
let mut invalid_utf8 = false;
for token in raw_tokens {
let token = unsafe { CStr::from_ptr(*token) };
match token.to_str() {
Ok(token) => tokens.push(token.to_owned()),
Err(_) => invalid_utf8 = true,
}
}
unsafe { adamo_sys::adamo_liveliness_tokens_free(raw, count) };
if invalid_utf8 {
return Err(Error::InvalidUtf8);
}
Ok(tokens)
}
pub fn on_liveliness<F>(
&self,
pattern: &str,
history: bool,
callback: F,
) -> Result<LivelinessSubscriber<'_>>
where
F: Fn(String, bool) + Send + Sync + 'static,
{
let pattern = CString::new(pattern)?;
let mut state = Box::new(LivelinessState {
callback: Box::new(callback),
});
let raw = unsafe {
adamo_sys::adamo_liveliness_subscribe(
self.raw.as_ptr(),
pattern.as_ptr(),
history as i32,
Some(liveliness_trampoline),
state.as_mut() as *mut LivelinessState as *mut c_void,
)
};
match NonNull::new(raw) {
Some(raw) => Ok(LivelinessSubscriber {
raw,
_state: state,
_session: PhantomData,
}),
None => Err(last_ffi_error()),
}
}
}
impl Drop for Session {
fn drop(&mut self) {
unsafe { adamo_sys::adamo_session_free(self.raw.as_ptr()) };
}
}
fn json_escape(s: &str) -> String {
let mut out = String::with_capacity(s.len());
for c in s.chars() {
match c {
'"' => out.push_str("\\\""),
'\\' => out.push_str("\\\\"),
'\n' => out.push_str("\\n"),
'\r' => out.push_str("\\r"),
'\t' => out.push_str("\\t"),
c if (c as u32) < 0x20 => {
use std::fmt::Write;
let _ = write!(out, "\\u{:04x}", c as u32);
}
c => out.push(c),
}
}
out
}
fn json_quote(s: &str) -> String {
let mut out = String::with_capacity(s.len() + 2);
out.push('"');
out.push_str(&json_escape(s));
out.push('"');
out
}
#[derive(Debug, Clone, Copy)]
pub struct PublishOptions {
pub priority: u8,
pub express: bool,
}
pub struct Priority;
impl Priority {
pub const REAL_TIME: u8 = 250;
pub const INTERACTIVE_HIGH: u8 = 220;
pub const INTERACTIVE_LOW: u8 = 190;
pub const DATA_HIGH: u8 = 150;
pub const DATA: u8 = 100;
pub const DATA_LOW: u8 = 80;
pub const BACKGROUND: u8 = 20;
}
impl Default for PublishOptions {
fn default() -> Self {
Self {
priority: Priority::DATA,
express: false,
}
}
}
#[derive(Debug, Clone, Copy)]
pub struct PublisherOptions {
pub priority: u8,
pub express: bool,
pub reliable: bool,
}
impl Default for PublisherOptions {
fn default() -> Self {
Self {
priority: Priority::DATA,
express: false,
reliable: false,
}
}
}
pub struct Publisher<'a> {
raw: NonNull<adamo_sys::adamo_publisher_t>,
_session: PhantomData<&'a Session>,
}
unsafe impl Send for Publisher<'_> {}
unsafe impl Sync for Publisher<'_> {}
impl Publisher<'_> {
pub fn put(&self, payload: &[u8]) -> Result<()> {
let rc = unsafe {
adamo_sys::adamo_publisher_put(self.raw.as_ptr(), payload.as_ptr(), payload.len())
};
if rc == 0 { Ok(()) } else { Err(last_ffi_error()) }
}
}
impl Drop for Publisher<'_> {
fn drop(&mut self) {
unsafe { adamo_sys::adamo_publisher_free(self.raw.as_ptr()) };
}
}
pub struct StateStore<'a> {
raw: NonNull<adamo_sys::adamo_state_t>,
_session: PhantomData<&'a Session>,
}
unsafe impl Send for StateStore<'_> {}
unsafe impl Sync for StateStore<'_> {}
impl StateStore<'_> {
pub fn set(&self, key: &str, value: &[u8]) -> Result<()> {
let key = CString::new(key)?;
let rc = unsafe {
adamo_sys::adamo_state_set(self.raw.as_ptr(), key.as_ptr(), value.as_ptr(), value.len())
};
if rc == 0 { Ok(()) } else { Err(last_ffi_error()) }
}
pub fn get(&self, key: &str) -> Result<Option<Vec<u8>>> {
let key = CString::new(key)?;
let mut ptr: *mut u8 = std::ptr::null_mut();
let mut len: usize = 0;
let rc =
unsafe { adamo_sys::adamo_state_get(self.raw.as_ptr(), key.as_ptr(), &mut ptr, &mut len) };
if rc != 0 {
return Err(last_ffi_error());
}
if ptr.is_null() || len == 0 {
return Ok(None);
}
let bytes = unsafe { slice::from_raw_parts(ptr, len) }.to_vec();
unsafe { adamo_sys::adamo_state_bytes_free(ptr, len) };
Ok(Some(bytes))
}
pub fn delete(&self, key: &str) -> Result<()> {
let key = CString::new(key)?;
let rc = unsafe { adamo_sys::adamo_state_delete(self.raw.as_ptr(), key.as_ptr()) };
if rc == 0 { Ok(()) } else { Err(last_ffi_error()) }
}
}
impl Drop for StateStore<'_> {
fn drop(&mut self) {
unsafe { adamo_sys::adamo_state_free(self.raw.as_ptr()) };
}
}
pub struct TaskRunner<'a> {
raw: NonNull<adamo_sys::adamo_task_runner_t>,
_session: PhantomData<&'a Session>,
}
unsafe impl Send for TaskRunner<'_> {}
unsafe impl Sync for TaskRunner<'_> {}
impl TaskRunner<'_> {
pub fn select(&self, task_set_id: &str, name: &str, subtasks: &[(&str, &str)]) -> Result<()> {
let task_set_id = CString::new(task_set_id)?;
let name = CString::new(name)?;
let ids: Vec<CString> = subtasks
.iter()
.map(|(id, _)| CString::new(*id))
.collect::<std::result::Result<_, _>>()?;
let names: Vec<CString> = subtasks
.iter()
.map(|(_, n)| CString::new(*n))
.collect::<std::result::Result<_, _>>()?;
let id_ptrs: Vec<*const std::os::raw::c_char> = ids.iter().map(|c| c.as_ptr()).collect();
let name_ptrs: Vec<*const std::os::raw::c_char> =
names.iter().map(|c| c.as_ptr()).collect();
let rc = unsafe {
adamo_sys::adamo_task_runner_select(
self.raw.as_ptr(),
task_set_id.as_ptr(),
name.as_ptr(),
id_ptrs.as_ptr(),
name_ptrs.as_ptr(),
subtasks.len(),
)
};
if rc == 0 { Ok(()) } else { Err(last_ffi_error()) }
}
pub fn advance(&self) -> Result<()> {
let rc = unsafe { adamo_sys::adamo_task_runner_advance(self.raw.as_ptr()) };
if rc == 0 { Ok(()) } else { Err(last_ffi_error()) }
}
pub fn reset(&self) -> Result<()> {
let rc = unsafe { adamo_sys::adamo_task_runner_reset(self.raw.as_ptr()) };
if rc == 0 { Ok(()) } else { Err(last_ffi_error()) }
}
}
impl Drop for TaskRunner<'_> {
fn drop(&mut self) {
unsafe { adamo_sys::adamo_task_runner_free(self.raw.as_ptr()) };
}
}
pub struct Subscriber<'a> {
raw: NonNull<adamo_sys::adamo_subscriber_t>,
_session: PhantomData<&'a Session>,
}
unsafe impl Send for Subscriber<'_> {}
impl Subscriber<'_> {
pub fn recv(&self, timeout: Option<Duration>) -> Result<Sample> {
let ms = timeout.map_or(0u64, |d| d.as_millis().min(u64::MAX as u128) as u64);
let raw = unsafe { adamo_sys::adamo_sub_recv(self.raw.as_ptr(), ms) };
Sample::from_raw(raw).ok_or_else(|| {
match last_ffi_error() {
Error::Ffi(m) if m == "(no error message)" => Error::Timeout,
e => e,
}
})
}
pub fn try_recv(&self) -> Result<Option<Sample>> {
let raw = unsafe { adamo_sys::adamo_sub_try_recv(self.raw.as_ptr()) };
if let Some(sample) = Sample::from_raw(raw) {
return Ok(Some(sample));
}
match last_ffi_error() {
Error::Ffi(m) if m == "(no error message)" => Ok(None),
e => Err(e),
}
}
}
impl Drop for Subscriber<'_> {
fn drop(&mut self) {
unsafe { adamo_sys::adamo_sub_free(self.raw.as_ptr()) };
}
}
type SampleCallback = dyn Fn(Sample) + Send + Sync + 'static;
struct CallbackState {
callback: Box<SampleCallback>,
}
unsafe extern "C" fn callback_trampoline(
sample: *const adamo_sys::adamo_sample_t,
user: *mut c_void,
) {
if sample.is_null() || user.is_null() {
return;
}
let Some(sample) = Sample::from_borrowed(sample) else {
return;
};
let state = unsafe { &*(user as *const CallbackState) };
if catch_unwind(AssertUnwindSafe(|| (state.callback)(sample))).is_err() {
eprintln!("adamo subscribe_with callback panicked");
}
}
pub struct CallbackSubscriber<'a> {
raw: NonNull<adamo_sys::adamo_cb_sub_t>,
_state: Box<CallbackState>,
_session: PhantomData<&'a Session>,
}
unsafe impl Send for CallbackSubscriber<'_> {}
unsafe impl Sync for CallbackSubscriber<'_> {}
impl Drop for CallbackSubscriber<'_> {
fn drop(&mut self) {
unsafe { adamo_sys::adamo_cb_sub_free(self.raw.as_ptr()) };
}
}
pub struct LivelinessToken<'a> {
raw: NonNull<adamo_sys::adamo_liveliness_token_t>,
_session: PhantomData<&'a Session>,
}
unsafe impl Send for LivelinessToken<'_> {}
unsafe impl Sync for LivelinessToken<'_> {}
impl Drop for LivelinessToken<'_> {
fn drop(&mut self) {
unsafe { adamo_sys::adamo_liveliness_token_free(self.raw.as_ptr()) };
}
}
type LivelinessCallback = dyn Fn(String, bool) + Send + Sync + 'static;
struct LivelinessState {
callback: Box<LivelinessCallback>,
}
unsafe extern "C" fn liveliness_trampoline(
key: *const c_char,
alive: i32,
user: *mut c_void,
) {
if key.is_null() || user.is_null() {
return;
}
let key = unsafe { CStr::from_ptr(key).to_string_lossy().into_owned() };
let state = unsafe { &*(user as *const LivelinessState) };
if catch_unwind(AssertUnwindSafe(|| (state.callback)(key, alive != 0))).is_err() {
eprintln!("adamo on_liveliness callback panicked");
}
}
pub struct LivelinessSubscriber<'a> {
raw: NonNull<adamo_sys::adamo_liveliness_sub_t>,
_state: Box<LivelinessState>,
_session: PhantomData<&'a Session>,
}
unsafe impl Send for LivelinessSubscriber<'_> {}
unsafe impl Sync for LivelinessSubscriber<'_> {}
impl Drop for LivelinessSubscriber<'_> {
fn drop(&mut self) {
unsafe { adamo_sys::adamo_liveliness_sub_free(self.raw.as_ptr()) };
}
}
#[derive(Debug, Clone)]
pub struct Sample {
pub key: String,
pub payload: Vec<u8>,
pub is_delete: bool,
pub timestamp_us: Option<u64>,
}
impl Sample {
fn from_raw(raw: *mut adamo_sys::adamo_sample_t) -> Option<Self> {
if raw.is_null() {
return None;
}
let sample = unsafe {
let s = &*raw;
let key = if s.key.is_null() {
String::new()
} else {
CStr::from_ptr(s.key).to_string_lossy().into_owned()
};
let payload = if s.payload.is_null() || s.payload_len == 0 {
Vec::new()
} else {
slice::from_raw_parts(s.payload, s.payload_len).to_vec()
};
let is_delete = s.is_delete != 0;
let timestamp_us = (s.timestamp_us != 0).then_some(s.timestamp_us);
adamo_sys::adamo_sample_free(raw);
Sample {
key,
payload,
is_delete,
timestamp_us,
}
};
Some(sample)
}
fn from_borrowed(raw: *const adamo_sys::adamo_sample_t) -> Option<Self> {
if raw.is_null() {
return None;
}
let s = unsafe { &*raw };
let key = if s.key.is_null() {
String::new()
} else {
unsafe { CStr::from_ptr(s.key).to_string_lossy().into_owned() }
};
let payload = if s.payload.is_null() || s.payload_len == 0 {
Vec::new()
} else {
unsafe { slice::from_raw_parts(s.payload, s.payload_len).to_vec() }
};
Some(Sample {
key,
payload,
is_delete: s.is_delete != 0,
timestamp_us: (s.timestamp_us != 0).then_some(s.timestamp_us),
})
}
}