|
@@ -16,9 +16,10 @@
|
|
|
* along with this program. If not, see <https://www.gnu.org/licenses/>.
|
|
* along with this program. If not, see <https://www.gnu.org/licenses/>.
|
|
|
*/
|
|
*/
|
|
|
|
|
|
|
|
|
|
+use async_lock::Mutex as AsyncMutex;
|
|
|
use darkfi::{
|
|
use darkfi::{
|
|
|
event_graph::{self},
|
|
event_graph::{self},
|
|
|
- net::transport::Dialer,
|
|
|
|
|
|
|
+ net::transport::{Dialer, PtStream},
|
|
|
system::ExecutorPtr,
|
|
system::ExecutorPtr,
|
|
|
util::path::expand_path,
|
|
util::path::expand_path,
|
|
|
Error, Result,
|
|
Error, Result,
|
|
@@ -27,10 +28,20 @@ use darkfi_serial::{
|
|
|
async_trait, deserialize_async_partial, AsyncDecodable, AsyncEncodable, Encodable,
|
|
async_trait, deserialize_async_partial, AsyncDecodable, AsyncEncodable, Encodable,
|
|
|
SerialDecodable, SerialEncodable,
|
|
SerialDecodable, SerialEncodable,
|
|
|
};
|
|
};
|
|
|
-use evgrd::{FetchEventsMessage, LocalEventGraph, VersionMessage, MSG_EVENT, MSG_FETCHEVENTS};
|
|
|
|
|
|
|
+use evgrd::{
|
|
|
|
|
+ FetchEventsMessage, LocalEventGraph, LocalEventGraphPtr, VersionMessage, MSG_EVENT,
|
|
|
|
|
+ MSG_FETCHEVENTS, MSG_SENDEVENT,
|
|
|
|
|
+};
|
|
|
use log::{error, info};
|
|
use log::{error, info};
|
|
|
use sled_overlay::sled;
|
|
use sled_overlay::sled;
|
|
|
-use smol::fs;
|
|
|
|
|
|
|
+use smol::{
|
|
|
|
|
+ fs,
|
|
|
|
|
+ io::{ReadHalf, WriteHalf},
|
|
|
|
|
+};
|
|
|
|
|
+use std::sync::{
|
|
|
|
|
+ atomic::{AtomicBool, Ordering},
|
|
|
|
|
+ Arc, Mutex as SyncMutex, Weak,
|
|
|
|
|
+};
|
|
|
use url::Url;
|
|
use url::Url;
|
|
|
|
|
|
|
|
use crate::scene::SceneNodePtr;
|
|
use crate::scene::SceneNodePtr;
|
|
@@ -47,79 +58,182 @@ pub struct Privmsg {
|
|
|
pub msg: String,
|
|
pub msg: String,
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
-pub async fn receive_msgs(sg_root: SceneNodePtr, ex: ExecutorPtr) -> Result<()> {
|
|
|
|
|
- let chatview_node = sg_root.lookup_node("/window/view/chatty").ok_or(Error::ConnectFailed)?;
|
|
|
|
|
|
|
+impl Privmsg {
|
|
|
|
|
+ pub fn new(channel: String, nick: String, msg: String) -> Self {
|
|
|
|
|
+ Self { channel, nick, msg }
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
|
|
|
- info!(target: "darkirc", "Instantiating DarkIRC event DAG");
|
|
|
|
|
- //let datastore = expand_path(EVGRDB_PATH)?;
|
|
|
|
|
- //fs::create_dir_all(&datastore).await?;
|
|
|
|
|
- let sled_db = sled::open(EVGRDB_PATH)?;
|
|
|
|
|
|
|
+pub type LocalDarkIRCPtr = Arc<LocalDarkIRC>;
|
|
|
|
|
|
|
|
- let evgr = LocalEventGraph::new(sled_db.clone(), "darkirc_dag", 1, ex.clone()).await?;
|
|
|
|
|
|
|
+pub struct LocalDarkIRC {
|
|
|
|
|
+ is_connected: AtomicBool,
|
|
|
|
|
+ /// The reading half of the transport stream
|
|
|
|
|
+ reader: AsyncMutex<Option<ReadHalf<Box<dyn PtStream>>>>,
|
|
|
|
|
+ /// The writing half of the transport stream
|
|
|
|
|
+ writer: AsyncMutex<Option<WriteHalf<Box<dyn PtStream>>>>,
|
|
|
|
|
|
|
|
- let endpoint = "tcp://127.0.0.1:5588";
|
|
|
|
|
- let endpoint = "tcp://192.168.1.38:5588";
|
|
|
|
|
- let endpoint = Url::parse(endpoint)?;
|
|
|
|
|
|
|
+ evgr: LocalEventGraphPtr,
|
|
|
|
|
+ receive_task: SyncMutex<Option<smol::Task<()>>>,
|
|
|
|
|
|
|
|
- let dialer = Dialer::new(endpoint.clone(), None).await?;
|
|
|
|
|
- let timeout = std::time::Duration::from_secs(60);
|
|
|
|
|
|
|
+ chatview_node: SceneNodePtr,
|
|
|
|
|
+}
|
|
|
|
|
|
|
|
- let mut stream = dialer.dial(Some(timeout)).await?;
|
|
|
|
|
- info!(target: "darkirc", "Connected to the backend: {endpoint}");
|
|
|
|
|
|
|
+impl LocalDarkIRC {
|
|
|
|
|
+ pub async fn new(sg_root: SceneNodePtr, ex: ExecutorPtr) -> Result<Arc<Self>> {
|
|
|
|
|
+ let chatview_node = sg_root.lookup_node("/window/view/chatty").unwrap();
|
|
|
|
|
|
|
|
- let version = VersionMessage::new();
|
|
|
|
|
- version.encode_async(&mut stream).await?;
|
|
|
|
|
|
|
+ info!(target: "darkirc", "Instantiating DarkIRC event DAG");
|
|
|
|
|
+ let datastore = expand_path(EVGRDB_PATH)?;
|
|
|
|
|
+ fs::create_dir_all(&datastore).await?;
|
|
|
|
|
+ let sled_db = sled::open(datastore)?;
|
|
|
|
|
+
|
|
|
|
|
+ let evgr = LocalEventGraph::new(sled_db.clone(), "darkirc_dag", 1, ex.clone()).await?;
|
|
|
|
|
+
|
|
|
|
|
+ Ok(Arc::new(Self {
|
|
|
|
|
+ is_connected: AtomicBool::new(false),
|
|
|
|
|
+ reader: AsyncMutex::new(None),
|
|
|
|
|
+ writer: AsyncMutex::new(None),
|
|
|
|
|
+
|
|
|
|
|
+ evgr,
|
|
|
|
|
+ receive_task: SyncMutex::new(None),
|
|
|
|
|
+
|
|
|
|
|
+ chatview_node,
|
|
|
|
|
+ }))
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ async fn reconnect(&self) -> Result<()> {
|
|
|
|
|
+ let endpoint = "tcp://127.0.0.1:5588";
|
|
|
|
|
+ let endpoint = Url::parse(endpoint)?;
|
|
|
|
|
+
|
|
|
|
|
+ let dialer = Dialer::new(endpoint.clone(), None).await?;
|
|
|
|
|
+ let timeout = std::time::Duration::from_secs(60);
|
|
|
|
|
+
|
|
|
|
|
+ let stream = dialer.dial(Some(timeout)).await?;
|
|
|
|
|
+ info!(target: "darkirc", "Connected to the backend: {endpoint}");
|
|
|
|
|
+
|
|
|
|
|
+ let (reader, writer) = smol::io::split(stream);
|
|
|
|
|
+ *self.writer.lock().await = Some(writer);
|
|
|
|
|
+ *self.reader.lock().await = Some(reader);
|
|
|
|
|
+
|
|
|
|
|
+ Ok(())
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ pub async fn start(self: Arc<Self>, ex: ExecutorPtr) -> Result<()> {
|
|
|
|
|
+ debug!(target: "darkirc", "LocalDarkIRC::start()");
|
|
|
|
|
+
|
|
|
|
|
+ self.version_exchange().await?;
|
|
|
|
|
+
|
|
|
|
|
+ let me = Arc::downgrade(&self);
|
|
|
|
|
+ let task = ex.spawn(async move {
|
|
|
|
|
+ while let Some(self_) = me.upgrade() {
|
|
|
|
|
+ self_.receive_msg().await.unwrap();
|
|
|
|
|
+ }
|
|
|
|
|
+ error!(target: "darkirc", "Closing DarkIRC receive loop");
|
|
|
|
|
+ });
|
|
|
|
|
|
|
|
- let server_version = VersionMessage::decode_async(&mut stream).await?;
|
|
|
|
|
- info!(target: "darkirc", "Backend server version: {}", server_version.protocol_version);
|
|
|
|
|
|
|
+ let mut receive_task = self.receive_task.lock().unwrap();
|
|
|
|
|
+ assert!(receive_task.is_none());
|
|
|
|
|
+ *receive_task = Some(task);
|
|
|
|
|
|
|
|
- let unref_tips = evgr.unreferenced_tips.read().await.clone();
|
|
|
|
|
- let fetchevs = FetchEventsMessage::new(unref_tips);
|
|
|
|
|
- MSG_FETCHEVENTS.encode_async(&mut stream).await?;
|
|
|
|
|
- fetchevs.encode_async(&mut stream).await?;
|
|
|
|
|
|
|
+ self.is_connected.store(true, Ordering::Relaxed);
|
|
|
|
|
|
|
|
- loop {
|
|
|
|
|
- let msg_type = u8::decode_async(&mut stream).await?;
|
|
|
|
|
|
|
+ Ok(())
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ async fn version_exchange(&self) -> Result<()> {
|
|
|
|
|
+ if !self.is_connected.load(Ordering::Relaxed) {
|
|
|
|
|
+ self.reconnect().await?;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ let mut writer = self.writer.lock().await;
|
|
|
|
|
+ let mut reader = self.reader.lock().await;
|
|
|
|
|
+ let writer = writer.as_mut().unwrap();
|
|
|
|
|
+ let reader = reader.as_mut().unwrap();
|
|
|
|
|
+
|
|
|
|
|
+ let version = VersionMessage::new();
|
|
|
|
|
+ version.encode_async(writer).await?;
|
|
|
|
|
+
|
|
|
|
|
+ let server_version = VersionMessage::decode_async(reader).await?;
|
|
|
|
|
+ info!(target: "darkirc", "Backend server version: {}", server_version.protocol_version);
|
|
|
|
|
+
|
|
|
|
|
+ let unref_tips = self.evgr.unreferenced_tips.read().await.clone();
|
|
|
|
|
+ let fetchevs = FetchEventsMessage::new(unref_tips);
|
|
|
|
|
+ MSG_FETCHEVENTS.encode_async(writer).await?;
|
|
|
|
|
+ fetchevs.encode_async(writer).await?;
|
|
|
|
|
+
|
|
|
|
|
+ Ok(())
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ async fn send_msg(&self, timestamp: u64, msg: Privmsg) -> Result<()> {
|
|
|
|
|
+ if !self.is_connected.load(Ordering::Relaxed) {
|
|
|
|
|
+ self.reconnect().await?;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ let mut writer = self.writer.lock().await;
|
|
|
|
|
+ let writer = writer.as_mut().unwrap();
|
|
|
|
|
+
|
|
|
|
|
+ MSG_SENDEVENT.encode_async(writer).await?;
|
|
|
|
|
+ timestamp.encode_async(writer).await?;
|
|
|
|
|
+ msg.encode_async(writer).await?;
|
|
|
|
|
+ Ok(())
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ async fn receive_msg(&self) -> Result<()> {
|
|
|
|
|
+ if !self.is_connected.load(Ordering::Relaxed) {
|
|
|
|
|
+ self.reconnect().await?;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ debug!(target: "darkirc", "Receiving message...");
|
|
|
|
|
+
|
|
|
|
|
+ let mut reader = self.reader.lock().await;
|
|
|
|
|
+ let reader = reader.as_mut().unwrap();
|
|
|
|
|
+
|
|
|
|
|
+ let msg_type = u8::decode_async(reader).await?;
|
|
|
debug!(target: "darkirc", "Received: {msg_type:?}");
|
|
debug!(target: "darkirc", "Received: {msg_type:?}");
|
|
|
if msg_type != MSG_EVENT {
|
|
if msg_type != MSG_EVENT {
|
|
|
error!(target: "darkirc", "Received invalid msg_type: {msg_type}");
|
|
error!(target: "darkirc", "Received invalid msg_type: {msg_type}");
|
|
|
- return Err(Error::MalformedPacket)
|
|
|
|
|
|
|
+ //return Err(Error::MalformedPacket)
|
|
|
|
|
+ return Ok(())
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- let ev = event_graph::Event::decode_async(&mut stream).await?;
|
|
|
|
|
|
|
+ let ev = event_graph::Event::decode_async(reader).await?;
|
|
|
|
|
|
|
|
- let genesis_timestamp = evgr.current_genesis.read().await.clone().timestamp;
|
|
|
|
|
|
|
+ let genesis_timestamp = self.evgr.current_genesis.read().await.clone().timestamp;
|
|
|
let ev_id = ev.id();
|
|
let ev_id = ev.id();
|
|
|
- if evgr.dag.contains_key(ev_id.as_bytes()).unwrap() ||
|
|
|
|
|
- !ev.validate(&evgr.dag, genesis_timestamp, evgr.days_rotation, None).await?
|
|
|
|
|
|
|
+ if self.evgr.dag.contains_key(ev_id.as_bytes()).unwrap() ||
|
|
|
|
|
+ !ev.validate(&self.evgr.dag, genesis_timestamp, self.evgr.days_rotation, None)
|
|
|
|
|
+ .await?
|
|
|
{
|
|
{
|
|
|
error!(target: "darkirc", "Event is invalid! {ev:?}");
|
|
error!(target: "darkirc", "Event is invalid! {ev:?}");
|
|
|
- continue
|
|
|
|
|
|
|
+ return Ok(())
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
debug!(target: "darkirc", "got {ev:?}");
|
|
debug!(target: "darkirc", "got {ev:?}");
|
|
|
- evgr.dag_insert(&[ev.clone()]).await.unwrap();
|
|
|
|
|
|
|
+ self.evgr.dag_insert(&[ev.clone()]).await.unwrap();
|
|
|
|
|
|
|
|
let privmsg: Privmsg = match deserialize_async_partial(ev.content()).await {
|
|
let privmsg: Privmsg = match deserialize_async_partial(ev.content()).await {
|
|
|
Ok((v, _)) => v,
|
|
Ok((v, _)) => v,
|
|
|
Err(e) => {
|
|
Err(e) => {
|
|
|
error!(target: "darkirc", "Failed deserializing incoming Privmsg event: {e}");
|
|
error!(target: "darkirc", "Failed deserializing incoming Privmsg event: {e}");
|
|
|
- continue
|
|
|
|
|
|
|
+ return Ok(())
|
|
|
}
|
|
}
|
|
|
};
|
|
};
|
|
|
|
|
|
|
|
debug!(target: "darkirc", "privmsg: {privmsg:?}");
|
|
debug!(target: "darkirc", "privmsg: {privmsg:?}");
|
|
|
|
|
|
|
|
if privmsg.channel != "#random" {
|
|
if privmsg.channel != "#random" {
|
|
|
- continue
|
|
|
|
|
|
|
+ return Ok(())
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
let mut arg_data = vec![];
|
|
let mut arg_data = vec![];
|
|
|
- ev.timestamp.encode(&mut arg_data).unwrap();
|
|
|
|
|
- ev.id().as_bytes().encode(&mut arg_data).unwrap();
|
|
|
|
|
- privmsg.nick.encode(&mut arg_data).unwrap();
|
|
|
|
|
- privmsg.msg.encode(&mut arg_data).unwrap();
|
|
|
|
|
|
|
+ ev.timestamp.encode_async(&mut arg_data).await.unwrap();
|
|
|
|
|
+ ev.id().as_bytes().encode_async(&mut arg_data).await.unwrap();
|
|
|
|
|
+ privmsg.nick.encode_async(&mut arg_data).await.unwrap();
|
|
|
|
|
+ privmsg.msg.encode_async(&mut arg_data).await.unwrap();
|
|
|
|
|
+
|
|
|
|
|
+ self.chatview_node.call_method("insert_line", arg_data).await.unwrap();
|
|
|
|
|
|
|
|
- chatview_node.call_method("insert_line", arg_data).await.unwrap();
|
|
|
|
|
|
|
+ Ok(())
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|