use std::{collections::HashMap, io::stdin};
use nt_client::{data::DataType, publish::GenericPublisher, Client, ClientHandle};
use tokio::{select, sync::{broadcast, mpsc}};
use tracing::Level;
#[tokio::main]
async fn main() {
tracing_subscriber::fmt()
.with_max_level(Level::INFO)
.init();
let (cancel_send, mut cancel_recv) = broadcast::channel(1);
let client = Client::new(Default::default());
select! {
_ = cancel_recv.recv() => std::process::exit(0),
res = client.connect_setup(|client| setup(client, cancel_send.clone())) => {
let _ = cancel_send.send(());
if let Err(err) = res {
eprintln!("{err}");
}
},
};
}
fn setup(client: &Client, cancel_send: broadcast::Sender<()>) {
let mut cancel_recv = cancel_send.subscribe();
let mut publishers = HashMap::new();
let handle = client.handle().clone();
let sub_topic = handle.topic("/tmp");
tokio::spawn(async move {
let sub_task = tokio::spawn(async move {
let mut sub = sub_topic.subscribe(Default::default()).await.unwrap();
loop {
let _ = sub.recv().await;
};
});
let (stdin_send, mut stdin_recv) = mpsc::channel(1);
let stdin_task = tokio::task::spawn_blocking(move || {
loop {
let mut command = String::new();
stdin().read_line(&mut command).unwrap();
if command.trim() == "quit" { break; };
let _ = stdin_send.blocking_send(command.trim().to_string());
};
});
let command_task = tokio::spawn(async move {
loop {
let Some(command) = stdin_recv.recv().await else { break; };
let mut segments = command.splitn(2, " ");
let Some(command) = segments.next() else {
eprintln!("malformed command");
continue;
};
let args = segments.next().unwrap_or("");
match command {
"help" => {
println!("NT 4.1 Publisher CLI");
println!("Commands:");
println!("- publish <type> <topic>");
println!(" publishes to <topic> with type <type>");
println!(" possible types: string, boolean, int");
println!();
println!("- set <id> <value>");
println!(" publish to <id> <value>");
println!();
println!("- unpublish <id>");
println!(" unpublishes <id>");
}
"publish" => match publish_command(args, &handle, &mut publishers).await {
Ok(id) => println!("publishing with id {id}"),
Err(CommandError(err)) => eprintln!("{err}"),
}
"set" => match set_command(args, &mut publishers).await {
Ok(()) => println!("successfully set topic"),
Err(CommandError(err)) => eprintln!("{err}"),
}
"unpublish" => match unpublish_command(args, &mut publishers).await {
Ok(()) => println!("successfully unpublished"),
Err(CommandError(err)) => eprintln!("{err}"),
}
_ => eprintln!("unknown command {command}"),
};
};
});
select! {
_ = cancel_recv.recv() => {},
_ = sub_task => {},
_ = stdin_task => {},
res = command_task => {
if let Err(err) = res {
eprintln!("{err}");
}
},
};
let _ = cancel_send.send(());
});
}
async fn publish_command(
args: &str,
handle: &ClientHandle,
publishers: &mut HashMap<i32, GenericPublisher>,
) -> Result<i32, CommandError> {
let mut args = args.split_whitespace();
let r#type = args.next().ok_or(CommandError("missing type arg".to_string()))?;
let topic_name = args.collect::<Vec<&str>>().join(" ");
if topic_name.is_empty() { return Err(CommandError("missing topic arg".to_string())); };
let r#type = match r#type {
"string" => DataType::String,
"boolean" => DataType::Boolean,
"int" => DataType::Int,
_ => return Err(CommandError(format!("unknown type {type}"))),
};
let topic = handle.topic(topic_name);
let publisher = topic.generic_publish(r#type, Default::default()).await?;
let id = publisher.id();
publishers.insert(id, publisher);
Ok(id)
}
async fn set_command(
args: &str,
publishers: &mut HashMap<i32, GenericPublisher>
) -> Result<(), CommandError> {
let mut args = args.split_whitespace();
let id = args.next()
.ok_or(CommandError("missing id arg".to_string()))?
.parse::<i32>()?;
let value = args.collect::<Vec<&str>>().join(" ");
if value.is_empty() { return Err(CommandError("missing topic arg".to_string())); };
let publisher = publishers.get(&id).ok_or(CommandError("no publisher found".to_string()))?;
match publisher.data_type() {
DataType::String => publisher.set(value).await?,
DataType::Boolean => publisher.set(value.parse::<bool>()?).await?,
DataType::Int => publisher.set(value.parse::<i64>()?).await?,
r#type => return Err(CommandError(format!("unsupported data type {type:?}"))),
};
Ok(())
}
async fn unpublish_command(
args: &str,
publishers: &mut HashMap<i32, GenericPublisher>
) -> Result<(), CommandError> {
let mut args = args.split_whitespace();
let id = args.next()
.ok_or(CommandError("missing id arg".to_string()))?
.parse::<i32>()?;
publishers.remove(&id).ok_or(CommandError("publisher not found".to_string()))?;
Ok(())
}
struct CommandError(String);
impl<T: ToString> From<T> for CommandError {
fn from(value: T) -> Self {
Self(value.to_string())
}
}