diff --git a/src/admin.rs b/src/admin.rs index 920aad0..1ebc429 100644 --- a/src/admin.rs +++ b/src/admin.rs @@ -1,6 +1,7 @@ use fluvio::FluvioAdmin as FluvioAdminNative; use fluvio_sc_schema::topic::TopicSpec; use fluvio_future::task::run_block_on; +use fluvio_sc_schema::objects::ListFilter; pub struct FluvioAdmin { pub inner: FluvioAdminNative } @@ -10,11 +11,29 @@ impl FluvioAdmin { run_block_on(FluvioAdminNative::connect()).map(|a| Box::new(FluvioAdmin { inner: a })).map_err(|e| e.to_string()) } + pub fn connect_with_config(config: &super::FluvioConfig) -> Result, String> { + run_block_on(FluvioAdminNative::connect_with_config(&config.inner)) + .map(|a| Box::new(FluvioAdmin { inner: a })) + .map_err(|e| e.to_string()) + } + + pub fn topic_exists(self: &Self, topic: &str) -> bool { + run_block_on(self.inner.list::(Vec::new())) + .map(|topics| topics.into_iter().any(|t| t.name == topic)) + .unwrap_or(false) + } + pub fn create_topic(self: &Self, topic: &str, partitions: i32, replicas: i32) -> Result<(), String> { run_block_on(self.inner.create(topic.to_string(), false, TopicSpec::new_computed(partitions as u32, replicas as u32, None))) .map_err(|e| e.to_string()) } + pub fn list_topics(&self) -> Result, String> { + run_block_on(self.inner.list::(Vec::new())) + .map(|topics| topics.into_iter().map(|t| t.name).collect()) + .map_err(|e| e.to_string()) + } + pub fn delete_topic(self: &Self, topic: &str) -> Result<(), String> { run_block_on(self.inner.delete::(topic.to_string())) .map_err(|e| e.to_string()) diff --git a/src/lib.rs b/src/lib.rs index 4d45b6d..2b1397f 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -118,6 +118,13 @@ mod ffi { /// Connects to the Fluvio Administrative controller #[Self = "FluvioAdmin"] fn connect() -> Result>; + /// Connects to the Fluvio Administrative controller with explicit config + #[Self = "FluvioAdmin"] + fn connect_with_config(config: &FluvioConfig) -> Result>; + /// Lists all topics + fn list_topics(self: &FluvioAdmin) -> Result>; + /// Checks if a topic exists + fn topic_exists(self: &FluvioAdmin, topic: &str) -> bool; /// Dispatches a command to create a new topic fn create_topic(self: &FluvioAdmin, topic: &str, partitions: i32, replicas: i32) -> Result<()>; /// Dispatches a command to violently delete a topic