#![deny(unsafe_code)]
#![allow(dead_code)]
use async_std::task;
use futures::{sink::SinkExt, stream::StreamExt};
use log::{error, info};
use koibumi_common::{boxes::Boxes, param::Params};
use koibumi_core::{address::Address, message};
use koibumi_node::{self as node, Command, Event, Response};
fn handle_msg(boxes: &Boxes, user_id: Vec<u8>, address: Address, object: message::Object) {
let identity = boxes.user().private_identity_by_address(&address);
if identity.is_none() {
error!("identity not found for address: {}", address);
return
}
let identity = identity.unwrap();
task::block_on(async {
match boxes.manager().insert_msg(user_id, identity, object).await {
Ok(message) => {
println!("From: {}", message.from_address().to_string());
println!("{}", String::from_utf8_lossy(message.content()).to_string());
}
Err(err) => {
error!("{}", err);
return;
}
}
});
}
fn handle_broadcast(boxes: &Boxes, user_id: Vec<u8>, address: Address, object: message::Object) {
task::block_on(async {
match boxes
.manager()
.insert_broadcast(user_id, address, object)
.await
{
Ok(message) => {
println!("From: {}", message.from_address().to_string());
println!("{}", String::from_utf8_lossy(message.content()).to_string());
}
Err(err) => {
error!("{}", err);
return;
}
}
});
}
fn main() {
let params = Params::new();
koibumi_common::log::init(¶ms).unwrap_or_else(|err| {
println!("Warning: Failed to initialize logger.");
println!("{}", err);
});
let config = koibumi_common::config::load(¶ms).unwrap_or_else(|err| {
error!("Failed to load config file: {}", err);
std::process::exit(1)
});
let mut boxes = match task::block_on(koibumi_common::boxes::prepare(¶ms)) {
Ok(boxes) => Some(boxes),
Err(err) => {
error!("{}", err);
None
}
};
let (command_sender, mut response_receiver, node_handle) = node::spawn();
info!("Start");
if boxes.is_none() {
error!("No boxes");
std::process::exit(1)
}
#[cfg(feature = "ctrlc")]
{
use std::sync::atomic::{AtomicUsize, Ordering};
use async_std::sync::Arc;
let sender = command_sender.clone();
let ctrlc_count = Arc::new(AtomicUsize::new(0));
let cc = Arc::clone(&ctrlc_count);
ctrlc::set_handler(move || match cc.fetch_add(1, Ordering::SeqCst) {
0 => {
let mut sender = sender.clone();
task::block_on(async {
sender.send(Command::Stop).await.unwrap_or_else(|err| {
error!("{}", err);
});
});
}
1 => {
let mut sender = sender.clone();
task::block_on(async {
sender.send(Command::Abort).await.unwrap_or_else(|err| {
error!("{}", err);
});
});
}
_ => std::process::exit(0),
})
.unwrap_or_else(|err| {
error!("{}", err);
});
}
let mut sender = command_sender;
let response = task::block_on(async {
let pool = koibumi_common::node::prepare(¶ms)
.await
.unwrap_or_else(|err| {
error!("{}", err);
std::process::exit(1);
});
let users = vec![boxes.as_ref().unwrap().user().clone().into()];
if let Err(err) = sender
.send(Command::Start(config.into(), pool, users))
.await
{
error!("{}", err);
return None;
}
response_receiver.next().await
});
if response.is_none() {
error!("Could not start node.");
std::process::exit(1)
}
let Response::Started(mut receiver) = response.unwrap();
task::block_on(async {
while let Some(event) = receiver.next().await {
match event {
Event::ConnectionCounts { .. } => (),
Event::AddrCount(_count) => (),
Event::Established {
addr,
user_agent,
rating,
} => {
info!("established: {} {} rating:{}", addr, user_agent, rating);
}
Event::Disconnected { addr } => {
info!("disconnected: {}", addr);
}
Event::Objects { .. } => (),
Event::Stopped => {
boxes = None;
std::process::exit(0);
}
Event::Msg {
user_id,
address,
object,
} => {
if let Some(boxes) = &boxes {
handle_msg(boxes, user_id, address, object);
}
}
Event::Broadcast {
user_id,
address,
object,
} => {
if let Some(boxes) = &boxes {
handle_broadcast(boxes, user_id, address, object);
}
}
}
}
node_handle.await;
});
}