/* This file is part of DarkFi (https://dark.fi)
*
* Copyright (C) 2020-2026 Dyne.org foundation
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU Affero General Public License as
* published by the Free Software Foundation, either version 3 of the
* License, or (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU Affero General Public License for more details.
*
* You should have received a copy of the GNU Affero General Public License
* along with this program. If not, see .
*/
use std::{
collections::{HashMap, HashSet},
io::Cursor,
sync::{Arc, OnceLock, Weak},
time::UNIX_EPOCH,
};
use async_lock::RwLock;
use crypto_box::{ChaChaBox, PublicKey, SecretKey};
use darkfi::{
event_graph::{
self,
proto::{EventPut, ProtocolEventGraph},
EventGraph, EventGraphConfig, EventGraphPtr,
},
net::{
dnet::DnetEvent,
session::SESSION_DEFAULT,
settings::{MagicBytes, NetworkProfile, Settings as NetSettings},
ChannelPtr, P2p, P2pPtr,
},
system::{sleep, Subscription},
Result as DarkFiResult,
};
use darkfi_serial::{
deserialize_async, serialize, serialize_async, AsyncEncodable, Decodable, Encodable,
};
use irc2::{
crypto::saltbox,
irc::{server::MAX_NICK_LEN, IrcChannel, IrcContact},
pad, unpad, Privmsg,
};
use parking_lot::Mutex as SyncMutex;
use sled_overlay::sled;
use crate::{
app::schema::menu::{channel::Channel, contact::Contact},
error::{Error, Result},
prop::{BatchGuardPtr, PropertyAtomicGuard, PropertyPtr, PropertyStr, Role},
scene::{MethodCallSub, Pimpl, SceneNode, SceneNodePtr, SceneNodeType, SceneNodeWeak, Slot},
ui::{
chatview::{MessageId, Timestamp},
OnModify,
},
ExecutorPtr,
};
use super::PluginSettings;
const P2P_RETRY_TIME: u64 = 20;
const COOLOFF_SLEEP_TIME: u64 = 20;
const COOLOFF_SYNC_ATTEMPTS: usize = 6;
const SYNC_MIN_PEERS: usize = 2;
pub(crate) const P2P_OUTBOUND_ACTIVE: usize = 6;
const P2P_OUTBOUND_SLEEP: usize = 1;
/// Update `outbound_peers` property useful for diagnostics
const DNET_ENABLED: bool = true;
/// Due to drift between different machine's clocks, if the message timestamp is recent
/// then we will just correct it to the current time so messages appear sequential in the UI.
const RECENT_TIME_DIST: u64 = 25_000;
// NOTE: if `paths` already lives in a shared module (e.g. `super::paths` from
// darkirc.rs's parent module), delete this block and add `use super::paths::*;`
// instead. Duplicated here so this file compiles standalone.
#[cfg(target_os = "android")]
mod paths {
use crate::android::{get_appdata_path, get_external_storage_path};
use std::path::PathBuf;
pub fn get_evgrdb_path() -> PathBuf {
get_external_storage_path().join("evgr2")
}
pub fn get_chatdb_path() -> PathBuf {
get_external_storage_path().join("chatdb")
}
pub fn get_use_tor_filename() -> PathBuf {
get_external_storage_path().join("use_tor.txt")
}
pub fn nick_filename() -> PathBuf {
get_appdata_path().join("/nick2.txt")
}
pub fn p2p_datastore_path() -> PathBuf {
get_appdata_path().join("darkirc2_p2p")
}
pub fn hostlist_path() -> PathBuf {
get_appdata_path().join("hostlist2.tsv")
}
}
#[cfg(not(target_os = "android"))]
mod paths {
use std::path::PathBuf;
pub fn get_evgrdb_path() -> PathBuf {
dirs::data_local_dir().unwrap().join("darkfi/app/evgr2")
}
pub fn get_chatdb_path() -> PathBuf {
dirs::data_local_dir().unwrap().join("darkfi/app/chatdb")
}
pub fn get_use_tor_filename() -> PathBuf {
dirs::data_local_dir().unwrap().join("darkfi/app/use_tor.txt")
}
pub fn nick_filename() -> PathBuf {
dirs::cache_dir().unwrap().join("darkfi/app/nick2.txt")
}
pub fn p2p_datastore_path() -> PathBuf {
dirs::cache_dir().unwrap().join("darkfi/app/darkirc2_p2p")
}
pub fn hostlist_path() -> PathBuf {
dirs::cache_dir().unwrap().join("darkfi/app/hostlist2.tsv")
}
}
use paths::*;
macro_rules! t { ($($arg:tt)*) => { trace!(target: "plugin::darkirc2", $($arg)*); } }
macro_rules! d { ($($arg:tt)*) => { debug!(target: "plugin::darkirc2", $($arg)*); } }
macro_rules! i { ($($arg:tt)*) => { info!(target: "plugin::darkirc2", $($arg)*); } }
macro_rules! e { ($($arg:tt)*) => { error!(target: "plugin::darkirc2", $($arg)*); } }
macro_rules! w { ($($arg:tt)*) => { warn!(target: "plugin::darkirc2", $($arg)*); } }
struct SeenMsg {
id: MessageId,
is_self: bool,
seen_times: usize,
}
struct SeenMessages {
seen: Vec,
}
impl SeenMessages {
fn new() -> Self {
Self { seen: vec![] }
}
fn get_status(&self, id: &MessageId) -> Option<&SeenMsg> {
self.seen.iter().find(|s| s.id == *id)
}
fn push(&mut self, id: MessageId, is_self: bool) {
self.seen.push(SeenMsg { id, is_self, seen_times: 0 });
}
}
pub type DarkIrcPtr = Arc;
pub struct DarkIrc {
node: SceneNodeWeak,
tasks: SyncMutex>>,
p2p: P2pPtr,
event_graph: EventGraphPtr,
seen_msgs: SyncMutex,
nick: PropertyStr,
pub channels: RwLock>,
pub contacts: RwLock>,
channels_tree: sled::Tree,
contacts_tree: sled::Tree,
dm_secret: SecretKey,
settings: PluginSettings,
ex: ExecutorPtr,
}
impl DarkIrc {
pub async fn new(
node: SceneNodeWeak,
sg_root: SceneNodePtr,
ex: ExecutorPtr,
db: sled::Db,
) -> Result {
let node_ref = &node.upgrade().unwrap();
let nick = PropertyStr::wrap(node_ref, Role::Internal, "nick", 0).unwrap();
let setting_root = Arc::new(SceneNode::new("setting", SceneNodeType::SettingRoot));
node_ref.link(setting_root.clone());
i!("Starting DarkIRC backend");
let evgr_path = get_evgrdb_path();
let evgr_db = match sled::open(&evgr_path) {
Ok(db) => db,
Err(err) => {
e!("Sled database '{}' failed to open: {err}!", evgr_path.display());
return Err(Error::SledDbErr)
}
};
let setting_tree = evgr_db.open_tree("settings")?;
// Use the unified db for reading channels (UI stores channels there)
let channels_tree = db.open_tree("channels")?;
i!("Opened channels tree from unified db");
let contacts_tree = db.open_tree("contacts")?;
i!("Opened contacts tree from unified db");
let dm_secret = Self::load_or_create_dm_identity(&db);
let dm_public_b58 = bs58::encode(dm_secret.public_key().to_bytes()).into_string();
// Expose our DM public key on the plugin node so it can be displayed/shared.
node_ref
.set_property_str(
&mut PropertyAtomicGuard::none(),
Role::Internal,
"dm_public",
&dm_public_b58,
)
.unwrap();
i!("DM identity public key (share with contacts): {dm_public_b58}");
let settings = PluginSettings { setting_root, sled_tree: setting_tree };
let mut p2p_settings: NetSettings = Default::default();
p2p_settings.magic_bytes = MagicBytes([251, 229, 199, 181]);
p2p_settings.app_version = semver::Version::parse("0.5.0").unwrap();
p2p_settings.app_name = "darkirc".to_string();
if get_use_tor_filename().exists() {
i!("Setup P2P network [tor]");
let mut tor_profile = NetworkProfile::tor_default();
tor_profile.outbound_connect_timeout = 60;
p2p_settings.profiles.insert("tor".to_string(), tor_profile);
p2p_settings.outbound_peer_discovery_cooloff_time = 60;
p2p_settings.seeds.push(
url::Url::parse(
"tor://g7fxelebievvpr27w7gt24lflptpw3jeeuvafovgliq5utdst6xyruyd.onion:25552",
)
.unwrap(),
);
p2p_settings.seeds.push(
url::Url::parse(
"tor://yvklzjnfmwxhyodhrkpomawjcdvcaushsj6torjz2gyd7e25f3gfunyd.onion:25552",
)
.unwrap(),
);
p2p_settings.active_profiles = vec!["tor".to_string()];
} else {
i!("Setup P2P network [clearnet]");
let mut profile = NetworkProfile::default();
profile.outbound_connect_timeout = 40;
profile.channel_handshake_timeout = 30;
p2p_settings.profiles.insert("tcp+tls".to_string(), profile);
p2p_settings.outbound_connections = 5;
p2p_settings.inbound_connections = 2;
p2p_settings.seeds.push(url::Url::parse("tcp+tls://lilith0.dark.fi:9600").unwrap());
p2p_settings.seeds.push(url::Url::parse("tcp+tls://lilith1.dark.fi:9600").unwrap());
p2p_settings.active_profiles = vec!["tcp+tls".to_string()];
}
p2p_settings.p2p_datastore = p2p_datastore_path().into_os_string().into_string().ok();
p2p_settings.hostlist = hostlist_path().into_os_string().into_string().ok();
settings.add_p2p_settings(&p2p_settings);
settings.load_settings();
settings.update_p2p_settings(&mut p2p_settings);
let p2p = match P2p::new(p2p_settings.clone(), ex.clone()).await {
Ok(p2p) => p2p,
Err(err) => {
e!("Create p2p network failed: {err}!");
return Err(Error::ServiceFailed)
}
};
if DNET_ENABLED {
i!("Enabling dnet outbound-slot event stream for outbound_peers property");
p2p.dnet_enable();
}
let event_graph = match EventGraph::new(
p2p.clone(),
db.clone(),
std::path::PathBuf::new(),
false,
EventGraphConfig {
initial_genesis: 1_704_067_200_000,
hours_rotation: 1,
genesis_contents: b"darkirc-v1".to_vec(),
rln_enabled: false,
pregenerated_identity_commitments: vec![],
max_dags: Some(24),
},
ex.clone(),
)
.await
{
Ok(evgr) => evgr,
Err(err) => {
e!("Create event graph failed: {err}!");
return Err(Error::ServiceFailed)
}
};
if let Ok(prev_nick) = std::fs::read_to_string(nick_filename()) {
nick.set(&mut PropertyAtomicGuard::none(), prev_nick);
}
let self_ = Arc::new(Self {
node: node.clone(),
tasks: SyncMutex::new(vec![]),
p2p,
event_graph,
seen_msgs: SyncMutex::new(SeenMessages::new()),
nick,
channels: RwLock::new(HashMap::new()),
contacts: RwLock::new(HashMap::new()),
channels_tree,
contacts_tree,
dm_secret,
settings,
ex: ex.clone(),
});
self_.load_channels_from_db().await;
self_.load_contacts_from_db().await;
self_.clone().start(sg_root, ex).await;
Ok(Pimpl::DarkIrc(self_))
}
async fn dag_sync(self: Arc, channel_sub: Subscription>) {
i!("Starting p2p network");
while let Err(err) = self.p2p.clone().start().await {
// This usually means we cannot listen on the inbound ports
e!("Failed to start p2p network: {err}!");
e!("Usually this means there is another process listening on the same ports.");
e!("Trying again in {P2P_RETRY_TIME} secs");
sleep(P2P_RETRY_TIME).await;
}
i!("Waiting for some P2P connections...");
let mut sync_attempt = 0;
// TODO: these should be configurable
let fast_mode = false;
let dags_count = 24;
loop {
if self.p2p.is_connected() {
let peers_count = self.p2p.peers_count();
self.notify_connect(peers_count, self.event_graph.is_synced()).await;
// Wait until we have enough connections
if peers_count < SYNC_MIN_PEERS {
i!("Connected to {peers_count} peers. Waiting for more connections.");
continue
}
i!("Got peer connection");
sync_attempt += 1;
// Cool off periodically
if sync_attempt > COOLOFF_SYNC_ATTEMPTS {
i!("Wasn't able to sync yet. Cooling off for {COOLOFF_SLEEP_TIME} then will try again.");
sleep(COOLOFF_SLEEP_TIME).await;
sync_attempt = 0;
}
i!("Syncing static DAG");
match self.event_graph.static_sync().await {
Ok(()) => {
i!("Static synced successfully");
// log_memory("after static sync");
}
Err(e) => {
e!("Failed syncing static graph: {e}");
self.p2p.stop().await;
break
}
}
i!("Syncing event DAG (attempt #{sync_attempt})");
// Sync mode is now per-call: full sync replays
// every event (heavy, used by archival nodes), fast
// sync only fetches headers (light, used by clients
// that don't need to re-verify history).
let sync_result = if fast_mode {
self.event_graph.sync_selected_headers(dags_count).await
} else {
self.event_graph.sync_selected(dags_count).await
};
match sync_result {
Ok(()) => {
i!(
"Event DAG synced successfully ({} mode, {} dag(s))",
if fast_mode { "fast" } else { "full" },
dags_count,
);
break
}
Err(e) => {
// TODO: Maybe at this point we should prune or something?
// TODO: Or maybe just tell the user to delete the DAG from FS.
e!("Failed syncing DAG ({e}), retrying...");
}
}
} else {
i!("Waiting for some P2P connections...");
sleep(COOLOFF_SLEEP_TIME).await;
}
}
let peers_count = self.p2p.peers_count();
self.notify_connect(peers_count, self.event_graph.is_synced()).await;
// Initial sync finished. Now just notify of connection changes
loop {
// Wait for a channel
if let Err(err) = channel_sub.receive().await {
w!("There was an error listening for channels. The service closed unexpectedly with error: {err}");
continue
}
let peers_count = self.p2p.peers_count();
self.notify_connect(peers_count, self.event_graph.is_synced()).await;
}
}
/// Send a notification when there's a change in number of peers or the DAG sync status
pub async fn notify_connect(&self, peers_count: usize, is_dag_synced: bool) {
let node = self.node.upgrade().unwrap();
node.trigger("connect", serialize(&(peers_count as u32, is_dag_synced))).await.unwrap();
}
/// Update the `outbound_peers` property with the outgoing connection slots addrs.
/// Allows us to monitor the network state of our p2p node.
async fn relay_outbound_slots(dnet_sub: Subscription, prop: PropertyPtr) {
loop {
let event = dnet_sub.receive().await;
let (slot, kind, addr) = match event {
DnetEvent::OutboundSlotConnected(info) => {
(info.slot, "connected", Some(info.addr.to_string()))
}
DnetEvent::OutboundSlotConnecting(info) => (info.slot, "connecting", None),
DnetEvent::OutboundSlotDisconnected(info) => (info.slot, "disconnected", None),
DnetEvent::OutboundSlotSleeping(info) => (info.slot, "sleeping", None),
_ => continue,
};
let mut atom = PropertyAtomicGuard::none();
let idx = slot as usize;
assert!(idx < prop.get_len());
match addr {
Some(addr) => prop.set_str(&mut atom, Role::Internal, idx, addr).unwrap(),
None => prop.set_null(&mut atom, Role::Internal, idx).unwrap(),
}
}
}
async fn relay_events(self: Arc, ev_sub: Subscription) {
loop {
let ev = ev_sub.receive().await;
// Try to deserialize the `Event`'s content into a `Privmsg`
let privmsg: Privmsg = match deserialize_async(ev.content()).await {
Ok(v) => v,
Err(e) => {
e!("[IRC CLIENT] Failed deserializing incoming Privmsg event: {}", e);
continue
}
};
// Route the message. An already-decrypted (plaintext) message names a
// channel we hold directly: encrypted channels arrive as base58
// ciphertext, so a channel key we recognise is plaintext by definition
// and is accepted as-is. Anything else must decrypt as a channel or DM;
// undecryptable traffic (base58 garbage in neither map) is silently dropped.
let mut privmsg = privmsg;
// Is this a plaintext channel?
let is_plaintext = self.channels.read().await.contains_key(&privmsg.channel);
if !is_plaintext && !self.try_decrypt(&mut privmsg, &self.nick.get()).await {
continue;
}
let mut timest = ev.header.timestamp;
let msg_id = msg_id(&privmsg, timest);
t!(
"Relaying ev_id={:?}, ev={ev:?}, msg_id={msg_id}, privmsg={privmsg:?}, timest={timest}",
ev.id(),
);
let is_self = {
let mut is_self = false;
let mut seen = self.seen_msgs.lock();
match seen.get_status(&msg_id) {
Some(msg) => {
is_self = msg.is_self;
if !msg.is_self || msg.seen_times > 1 {
w!("Skipping duplicate seen message: {msg_id}");
continue
}
}
None => {
seen.push(msg_id.clone(), false);
}
}
is_self
};
// This is a hack to make messages appear sequentially in the UI
let now_timest = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
if !is_self && timest.abs_diff(now_timest) < RECENT_TIME_DIST {
d!("Applied timestamp correction: <{timest}> => <{now_timest}>");
timest = now_timest;
}
// Workaround for the chatview hack. This nick is off limits!
let mut nick = privmsg.nick;
if nick == "NOTICE" {
nick = "noticer".to_string();
}
self.notify_recv(privmsg.channel, timest, msg_id, nick, privmsg.msg).await;
}
}
/// Send a notification about a new received message
pub async fn notify_recv(
&self,
channel: String,
timestamp: Timestamp,
id: MessageId,
nick: String,
msg: String,
) {
assert!(
channel.starts_with('#') || channel.starts_with('@'),
"notify_recv channel must be a \"#name\" channel or \"@name\" DM, got: {channel}"
);
let mut arg_data = vec![];
channel.encode(&mut arg_data).unwrap();
timestamp.encode(&mut arg_data).unwrap();
id.encode(&mut arg_data).unwrap();
nick.encode(&mut arg_data).unwrap();
msg.encode(&mut arg_data).unwrap();
let node = self.node.upgrade().unwrap();
node.trigger("recv", arg_data).await.unwrap();
}
async fn process_send(me: &Weak, sub: &MethodCallSub) -> bool {
let Ok(method_call) = sub.receive().await else {
d!("Event relayer closed");
return false
};
t!("method called: send({method_call:?})");
assert!(method_call.send_res.is_none());
fn decode_data(data: &[u8]) -> std::io::Result<(Timestamp, String, String)> {
let mut cur = Cursor::new(&data);
let timest = Timestamp::decode(&mut cur).unwrap();
let channel = String::decode(&mut cur)?;
let msg = String::decode(&mut cur)?;
Ok((timest, channel, msg))
}
let Ok((timest, channel, msg)) = decode_data(&method_call.data) else {
e!("send() method invalid arg data");
return true
};
let Some(self_) = me.upgrade() else {
// Should not happen
panic!("self destroyed before send_method_task was stopped!");
};
self_.handle_send(timest, channel, msg).await;
true
}
/// User wants to send a msg
async fn handle_send(&self, timest: Timestamp, channel: String, msg: String) {
let nick = self.nick.get();
// Send text to channel
d!("Sending privmsg: {timest} {channel}: <{nick}> {msg}");
let mut msg = Privmsg { version: 0, msg_type: 0, channel, nick, msg };
// DM layers use the "@name" UI id; strip it to the bare contact key and
// require the contact to exist, else refuse to broadcast.
if let Some(bare) = msg.channel.strip_prefix('@').map(str::to_string) {
msg.channel = bare;
if self.try_encrypt_dm(&mut msg).await.is_err() {
e!("Refusing to send DM to unknown contact");
return;
}
} else {
assert!(msg.channel.starts_with('#'), "channel name must start with #");
self.try_encrypt_channel(&mut msg).await;
}
let evgr = self.event_graph.clone();
let event = event_graph::Event::with_timestamp(timest, serialize_async(&msg).await, &evgr)
.await
.unwrap();
let msg_id = msg_id(&msg, timest);
// Keep track of our own messages so we don't apply timestamp correction to them
// which messes up the msg id.
{
let mut seen = self.seen_msgs.lock();
seen.push(msg_id.clone(), true);
}
// Broadcast the msg
let current_genesis = self.event_graph.current_genesis.read().await;
let dag_name = current_genesis.header.timestamp.to_string();
if let Err(e) = evgr.insert_signal_with_blob(&event, &[], &dag_name).await {
e!("Failed inserting new event to DAG: {}", e);
}
if let Err(e) = self.p2p.broadcast(&EventPut(event, vec![])).await {
e!("Event broadcast was not admitted: {e}");
}
}
/// Load channels from UI database and populate encryption keys
pub async fn load_channels_from_db(&self) {
let mut channels = self.channels.write().await;
for item in self.channels_tree.iter() {
let (key, val) = item.unwrap();
let channel_name = String::from_utf8_lossy(&key).to_string();
let ui_channel = deserialize_async::(&val).await.unwrap();
// Convert to IrcChannel with encryption
let full_name = format!("#{}", channel_name);
let mut irc_channel =
IrcChannel { topic: String::new(), nicks: HashSet::new(), saltbox: None };
if let Some(secret) = ui_channel.secret {
// Convert secret array to SecretKey first, then derive PublicKey
let secret_key = SecretKey::from_bytes(secret);
let public = secret_key.public_key();
let saltbox = ChaChaBox::new(&public, &secret_key);
// Log the secret in base58 for debugging
let secret_b58 = bs58::encode(secret).into_string();
irc_channel.saltbox = Some(Arc::new(saltbox));
}
let is_encrypted = irc_channel.saltbox.is_some();
channels.insert(full_name, irc_channel);
i!("Loaded channel: #{} (encrypted: {})", channel_name, is_encrypted);
}
}
/// Load (or generate on first run) the single global DM identity key.
fn load_or_create_dm_identity(db: &sled::Db) -> SecretKey {
let tree = db.open_tree("dm_identity").expect("cannot open dm_identity tree");
if let Ok(Some(stored)) = tree.get(b"secret") {
if stored.len() == 32 {
let arr: [u8; 32] = stored.as_ref().try_into().unwrap();
return SecretKey::from_bytes(arr);
}
}
let bytes: [u8; 32] = rand::random();
let _ = tree.insert(b"secret", bytes.to_vec());
let _ = tree.flush();
SecretKey::from_bytes(bytes)
}
/// Load contacts from the UI database and build their encryption boxes.
pub async fn load_contacts_from_db(&self) {
let mut contacts = self.contacts.write().await;
contacts.clear();
for item in self.contacts_tree.iter() {
let (key, val) = item.unwrap();
let name = String::from_utf8_lossy(&key).to_string();
let contact = deserialize_async::(&val).await.unwrap();
let their_public = PublicKey::from(contact.public);
let saltbox = Arc::new(ChaChaBox::new(&their_public, &self.dm_secret));
let self_saltbox =
Arc::new(ChaChaBox::new(&self.dm_secret.public_key(), &self.dm_secret));
contacts.insert(name.clone(), IrcContact { saltbox, self_saltbox });
i!("Loaded contact: {name}");
}
}
async fn rescan_channel_history(self: Arc, channel: String) {
i!("Starting background rescan for channel: {channel}");
// Fetch and order all events from the DAG (like darkirc does)
let Ok(dag_events) = self.event_graph.order_events().await else {
e!("Failed to fetch events from DAG");
return;
};
let mut found_count = 0;
for event in dag_events.iter() {
// Deserialize Privmsg
let mut privmsg = match deserialize_async::(event.content()).await {
Ok(pm) => pm,
Err(e) => {
t!("Not a Privmsg event, skipping");
continue;
}
};
// Try to decrypt (handles encrypted channels)
self.try_decrypt(&mut privmsg, &self.nick.get()).await;
// Check if message belongs to target channel
if privmsg.channel != channel {
continue;
}
found_count += 1;
// Calculate message ID
let timest = event.header.timestamp;
let msg_id = msg_id(&privmsg, timest);
// Send to ChatView via notify_recv (handles DB storage and duplicates)
self.notify_recv(
channel.clone(),
timest,
msg_id,
privmsg.nick.clone(),
privmsg.msg.clone(),
)
.await;
}
i!("Rescan complete for {channel}: found {found_count} messages");
}
async fn process_rescan(me: &Weak, sub: &MethodCallSub) -> bool {
let Ok(method_call) = sub.receive().await else {
d!("Rescan method closed");
return false
};
t!("method called: rescan({method_call:?})");
let Some(self_) = me.upgrade() else {
e!("DarkIrc destroyed before rescan completed");
return false
};
// Decode channel name from method data
let mut cur = std::io::Cursor::new(&method_call.data);
let Ok(channel) = String::decode(&mut cur) else {
e!("Rescan method called with invalid channel data");
return false
};
self_.load_channels_from_db().await;
self_.load_contacts_from_db().await;
let task = self_.ex.clone().spawn(self_.clone().rescan_channel_history(channel));
self_.tasks.lock().push(task);
true
}
async fn apply_settings(self_: Arc, _: BatchGuardPtr) {
self_.settings.save_settings();
let p2p_settings = self_.p2p.settings();
let mut write_guard = p2p_settings.write().await;
self_.settings.update_p2p_settings(&mut write_guard);
}
async fn process_reconnect(me: &Weak, sub: &MethodCallSub) -> bool {
let Ok(method_call) = sub.receive().await else {
d!("Reconnect method closed");
return false
};
t!("method called: reconnect({method_call:?})");
let Some(self_) = me.upgrade() else {
e!("DarkIrc destroyed before reconnect completed");
return false
};
self_.handle_reconnect().await;
true
}
/// User requested to reconnect
async fn handle_reconnect(&self) {
i!("Manual P2P reconnection triggered");
self.p2p.clone().stop().await;
while let Err(err) = self.p2p.clone().start().await {
e!("Failed to start P2P network: {err}!");
e!("Retrying in {P2P_RETRY_TIME} secs");
sleep(P2P_RETRY_TIME).await;
}
i!("P2P reconnection completed");
}
async fn start(self: Arc, sg_root: SceneNodePtr, ex: ExecutorPtr) {
i!("Registering EventGraph P2P protocol");
let event_graph_ = Arc::clone(&self.event_graph);
let registry = self.p2p.protocol_registry();
registry
.register(SESSION_DEFAULT, move |channel, _| {
let event_graph_ = event_graph_.clone();
async move { ProtocolEventGraph::init(event_graph_, channel).await.unwrap() }
})
.await;
let me = Arc::downgrade(&self);
let node = &self.node.upgrade().unwrap();
let method_sub = node.subscribe_method_call("send").unwrap();
let me2 = me.clone();
let send_method_task =
ex.spawn(async move { while Self::process_send(&me2, &method_sub).await {} });
let reconnect_method_sub = node.subscribe_method_call("reconnect").unwrap();
let me2 = me.clone();
let reconnect_method_task =
ex.spawn(
async move { while Self::process_reconnect(&me2, &reconnect_method_sub).await {} },
);
let rescan_method_sub = node.subscribe_method_call("rescan").unwrap();
let me2 = me.clone();
let rescan_method_task =
ex.spawn(async move { while Self::process_rescan(&me2, &rescan_method_sub).await {} });
let mut on_modify = OnModify::new(ex.clone(), self.node.clone(), me.clone());
async fn save_nick(self_: Arc, _batch: BatchGuardPtr) {
let _ = std::fs::write(nick_filename(), self_.nick.get());
}
on_modify.when_change(self.nick.prop(), save_nick);
// `apply_settings` is triggered if any setting changes
for setting_node in self.settings.setting_root.get_children().iter() {
on_modify.when_change(
setting_node.get_property("value").clone().unwrap(),
Self::apply_settings,
);
}
let ev_sub = self.event_graph.event_subscribe().await;
let ev_task = ex.spawn(self.clone().relay_events(ev_sub));
// Sync the DAG / check sync status
let channel_sub = self.p2p.hosts().subscribe_channel().await;
let dag_task = ex.spawn(self.clone().dag_sync(channel_sub));
// Subscribe to window start/stop signals for dynamic outbound connections
let window_node = sg_root.lookup_node("/window").unwrap();
let (start_slot, start_recv) = Slot::new("app_start");
window_node.register("start", start_slot).unwrap();
let p2p = self.p2p.clone();
let start_task = ex.spawn(async move {
while let Ok(_) = start_recv.recv().await {
i!("App started: set outbound connections to {P2P_OUTBOUND_ACTIVE}");
p2p.settings().write().await.outbound_connections = P2P_OUTBOUND_ACTIVE;
p2p.clone().reload().await;
}
});
let (stop_slot, stop_recv) = Slot::new("app_stop");
window_node.register("stop", stop_slot).unwrap();
let p2p = self.p2p.clone();
let stop_task = ex.spawn(async move {
while let Ok(_) = stop_recv.recv().await {
i!("App stopped: set outbound connections to {P2P_OUTBOUND_SLEEP}");
p2p.settings().write().await.outbound_connections = P2P_OUTBOUND_SLEEP;
p2p.clone().reload().await;
}
});
let mut tasks = vec![
send_method_task,
reconnect_method_task,
rescan_method_task,
ev_task,
dag_task,
start_task,
stop_task,
];
if DNET_ENABLED {
let dnet_sub = self.p2p.dnet_subscribe().await;
let node = self.node.upgrade().unwrap();
let prop = node.get_property("outbound_peers").unwrap();
let dnet_task = ex.spawn(Self::relay_outbound_slots(dnet_sub, prop));
tasks.push(dnet_task);
}
tasks.append(&mut on_modify.tasks);
*self.tasks.lock() = tasks;
}
/// Encrypt a channel `Privmsg` in place if the channel has a shared key.
/// Open channels with no key are left plaintext.
pub async fn try_encrypt_channel(&self, privmsg: &mut Privmsg) {
let guard = self.channels.read().await;
let Some((name, channel)) = guard.get_key_value(&privmsg.channel) else {
return;
};
let Some(saltbox) = &channel.saltbox else {
return;
};
privmsg.channel = saltbox::encrypt(saltbox, &[0x00; MAX_NICK_LEN]);
privmsg.nick = saltbox::encrypt(saltbox, &pad(&privmsg.nick));
privmsg.msg = saltbox::encrypt(saltbox, privmsg.msg.as_bytes());
d!("Successfully encrypted message for {name}");
}
/// Encrypt a DM `Privmsg` in place for the contact named by `privmsg.channel`
/// (the bare key, with no leading "@"). Fails if the contact is unknown so
/// the caller can refuse to broadcast a message no one could decrypt.
pub async fn try_encrypt_dm(&self, privmsg: &mut Privmsg) -> Result<()> {
let guard = self.contacts.read().await;
let Some((name, contact)) = guard.get_key_value(&privmsg.channel) else {
return Err(Error::ContactNotFound);
};
privmsg.channel = saltbox::encrypt(&contact.saltbox, &[0x00; MAX_NICK_LEN]);
privmsg.nick = saltbox::encrypt(&contact.self_saltbox, &[0x00; MAX_NICK_LEN]);
privmsg.msg = saltbox::encrypt(&contact.saltbox, privmsg.msg.as_bytes());
d!("Successfully encrypted DM for {name}");
Ok(())
}
/// Try decrypting a `Privmsg` as a channel message in place. Returns true on
/// success. Plaintext messages for a known keyless channel are accepted
/// as-is; everything else returns false.
pub async fn try_decrypt_channel(&self, privmsg: &mut Privmsg) -> bool {
let Ok(channel_ciphertext) = bs58::decode(&privmsg.channel).into_vec() else {
// Not encrypted: accept only if it names a channel we hold.
return self.channels.read().await.contains_key(&privmsg.channel);
};
let Ok(nick_ciphertext) = bs58::decode(&privmsg.nick).into_vec() else { return false };
let Ok(msg_ciphertext) = bs58::decode(&privmsg.msg).into_vec() else { return false };
for (name, channel) in self.channels.read().await.iter() {
let Some(saltbox) = &channel.saltbox else { continue };
if saltbox::try_decrypt(saltbox, &channel_ciphertext).is_none() {
continue
};
let Some(mut nick_dec) = saltbox::try_decrypt(saltbox, &nick_ciphertext) else {
w!("Could not decrypt nick ciphertext for channel: {name}");
continue
};
let Some(msg_dec) = saltbox::try_decrypt(saltbox, &msg_ciphertext) else {
w!("Could not decrypt message ciphertext for channel: {name}");
continue
};
unpad(&mut nick_dec);
privmsg.channel = name.to_string();
privmsg.nick = String::from_utf8_lossy(&nick_dec).into();
privmsg.msg = String::from_utf8_lossy(&msg_dec).into();
d!("Successfully decrypted message for {name}");
return true
}
false
}
/// Try decrypting a `Privmsg` as a DM in place. Returns true on success.
pub async fn try_decrypt_contact(&self, privmsg: &mut Privmsg, self_nickname: &str) -> bool {
let Ok(channel_ciphertext) = bs58::decode(&privmsg.channel).into_vec() else {
return false
};
let Ok(nick_ciphertext) = bs58::decode(&privmsg.nick).into_vec() else { return false };
let Ok(msg_ciphertext) = bs58::decode(&privmsg.msg).into_vec() else { return false };
for (name, contact) in self.contacts.read().await.iter() {
if saltbox::try_decrypt(&contact.saltbox, &channel_ciphertext).is_none() {
continue
};
let nick = if saltbox::try_decrypt(&contact.self_saltbox, &nick_ciphertext).is_some() {
String::from(self_nickname)
} else {
name.to_string()
};
let Some(msg_dec) = saltbox::try_decrypt(&contact.saltbox, &msg_ciphertext) else {
w!("Could not decrypt message ciphertext for contact: {name}");
continue
};
privmsg.channel = format!("@{}", name);
privmsg.nick = nick;
privmsg.msg = String::from_utf8_lossy(&msg_dec).into();
return true
}
false
}
/// Try decrypting a given potentially encrypted `Privmsg` object as a channel
/// and then as a DM. Returns true if either succeeded.
pub async fn try_decrypt(&self, privmsg: &mut Privmsg, self_nickname: &str) -> bool {
self.try_decrypt_channel(privmsg).await ||
self.try_decrypt_contact(privmsg, self_nickname).await
}
}
pub fn msg_id(privmsg: &Privmsg, timest: u64) -> MessageId {
let mut hasher = blake3::Hasher::new();
0u8.encode(&mut hasher).unwrap();
0u8.encode(&mut hasher).unwrap();
timest.encode(&mut hasher).unwrap();
privmsg.channel.encode(&mut hasher).unwrap();
privmsg.nick.encode(&mut hasher).unwrap();
privmsg.msg.encode(&mut hasher).unwrap();
MessageId(hasher.finalize().into())
}