Current section

Files

Jump to
fluvio native fluvio_ex src lib.rs
Raw

native/fluvio_ex/src/lib.rs

use async_std::future;
use async_std::task;
use fluvio::dataplane::record::ConsumerRecord;
use fluvio::dataplane::ErrorCode;
use fluvio::{Fluvio, Offset, PartitionConsumer, TopicProducer};
use futures_util::stream::BoxStream;
use futures_util::StreamExt;
use rustler::{Atom, Encoder, Env, Error, NifResult, NifTuple, ResourceArc, Term};
use std::sync::Mutex;
use std::time::Duration;
// TODO review connection and configs
pub mod atom;
pub struct FluvioResource {
pub fluvio: Mutex<Fluvio>,
}
#[derive(NifTuple)]
pub struct FluvioResourceResponse {
pub ok: rustler::Atom,
pub resource: ResourceArc<FluvioResource>,
}
pub struct ConsumerResource {
pub consumer: Mutex<PartitionConsumer>,
pub stream: Mutex<BoxStream<'static, Result<ConsumerRecord, ErrorCode>>>,
}
#[derive(NifTuple)]
pub struct ConsumerResourceResponse {
pub ok: rustler::Atom,
pub resource: ResourceArc<ConsumerResource>,
}
pub struct ProducerResource {
pub producer: Mutex<TopicProducer>,
}
#[derive(NifTuple)]
pub struct ProducerResourceResponse {
pub ok: rustler::Atom,
pub producer: ResourceArc<ProducerResource>,
}
#[rustler::nif]
fn next<'a>(
env: Env<'a>,
resource: ResourceArc<ConsumerResource>,
secs_timeout: u64,
) -> NifResult<Term<'a>> {
let mut stream = resource.stream.lock().unwrap();
let next_val = match task::block_on(future::timeout(
Duration::from_secs(secs_timeout),
stream.next(),
)) {
Ok(next_val) => next_val,
Err(err) => return Err(Error::Term(Box::new(err.to_string()))),
};
match next_val {
Some(Ok(record)) => Ok((
record.offset(),
record.partition(),
record.key(),
String::from_utf8_lossy(record.value()).to_string(),
record.timestamp(),
)
.encode(env)),
Some(Err(err)) => Err(Error::Term(Box::new(err.to_string()))),
None => Err(Error::Term(Box::new("returned nothing"))),
}
}
#[rustler::nif]
fn new_consumer(
fluvio_res: ResourceArc<FluvioResource>,
topic: String,
partition: i32,
offset_type: String,
offset_value: u32,
) -> NifResult<ConsumerResourceResponse> {
let offset = match offset_type.as_str() {
"from_beginning" => Offset::from_beginning(offset_value),
"from_end" => Offset::from_end(offset_value),
_ => {
return Err(Error::Term(Box::new(std::format!(
"Unsupported offset_type {:?}",
offset_type
))))
}
};
let fluvio = fluvio_res.fluvio.lock().unwrap();
let consumer = match task::block_on(fluvio.partition_consumer(topic, partition)) {
Ok(consumer) => consumer,
Err(err) => return Err(Error::Term(Box::new(err.to_string()))),
};
let stream = match task::block_on(consumer.stream(offset)) {
Ok(stream) => stream,
Err(err) => return Err(Error::Term(Box::new(err.to_string()))),
};
Ok(ConsumerResourceResponse {
ok: atom::ok(),
resource: ResourceArc::new(ConsumerResource {
consumer: Mutex::new(consumer),
stream: Mutex::new(Box::pin(stream)),
}),
})
}
#[rustler::nif]
fn connect() -> NifResult<FluvioResourceResponse> {
match task::block_on(Fluvio::connect()) {
Ok(fluvio) => Ok(FluvioResourceResponse {
ok: atom::ok(),
resource: ResourceArc::new(FluvioResource {
fluvio: Mutex::new(fluvio),
}),
}),
Err(err) => Err(Error::Term(Box::new(err.to_string()))),
}
}
#[rustler::nif]
fn new_producer(
fluvio_res: ResourceArc<FluvioResource>,
topic: String,
) -> NifResult<ProducerResourceResponse> {
let fluvio = fluvio_res.fluvio.lock().unwrap();
match task::block_on(fluvio.topic_producer(topic)) {
Ok(producer) => Ok(ProducerResourceResponse {
ok: atom::ok(),
producer: ResourceArc::new(ProducerResource {
producer: Mutex::new(producer),
}),
}),
Err(err) => Err(Error::Term(Box::new(err.to_string()))),
}
}
#[rustler::nif]
fn send(resource: ResourceArc<ProducerResource>, key: String, value: String) -> NifResult<Atom> {
let producer = resource.producer.lock().unwrap();
match task::block_on(producer.send(key, value)) {
Ok(_) => Ok(atom::ok()),
Err(err) => Err(Error::Term(Box::new(err.to_string()))),
}
}
#[rustler::nif]
fn flush(resource: ResourceArc<ProducerResource>) -> NifResult<Atom> {
let producer = resource.producer.lock().unwrap();
match task::block_on(producer.flush()) {
Ok(_) => Ok(atom::ok()),
Err(err) => Err(Error::Term(Box::new(err.to_string()))),
}
}
rustler::init!(
"Elixir.Fluvio.Native",
[new_producer, send, flush, new_consumer, next, connect],
load = on_load
);
fn on_load(env: Env, _info: Term) -> bool {
rustler::resource!(ProducerResource, env);
rustler::resource!(ConsumerResource, env);
rustler::resource!(FluvioResource, env);
true
}