/* 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 darkfi::{
net::{
session::{SESSION_DIRECT, SESSION_INBOUND},
settings::{MagicBytes, NetworkProfile, Settings as NetSettings},
P2p, P2pPtr,
},
system::{sleep, Publisher, PublisherPtr},
};
use darkfi_serial::{Decodable, Encodable};
use fud::{
event::FudEvent,
proto::ProtocolFud,
resource::ResourceStatus,
settings::Args as FudSettings,
util::{hash_to_string, FileSelection},
Fud,
};
use kvdb_overlay::Database;
use smol::lock::Mutex;
use std::{
collections::HashSet,
io::Cursor,
path::PathBuf,
sync::{Arc, OnceLock, Weak},
};
use url::Url;
use crate::{
error::{Error, Result},
prop::{PropertyAtomicGuard, PropertyBool, Role},
scene::{MethodCall, MethodCallSub, Pimpl, SceneNodePtr, SceneNodeWeak},
ui::chatview::FileMessageStatus,
ExecutorPtr,
};
const P2P_RETRY_TIME: u64 = 20;
#[cfg(target_os = "android")]
mod paths {
use crate::android::{get_appdata_path, get_external_storage_path};
use std::path::PathBuf;
pub fn get_base_path() -> PathBuf {
get_external_storage_path().join("fud")
}
pub fn get_db_path() -> PathBuf {
get_external_storage_path().join("fud/db")
}
pub fn get_downloads_path() -> PathBuf {
get_external_storage_path().join("fud/downloads")
}
pub fn get_use_tor_filename() -> PathBuf {
get_external_storage_path().join("use_tor.txt")
}
pub fn p2p_datastore_path() -> PathBuf {
get_appdata_path().join("fud/p2p")
}
pub fn hostlist_path() -> PathBuf {
get_appdata_path().join("fud/hostlist.tsv")
}
}
#[cfg(not(target_os = "android"))]
mod paths {
use std::path::PathBuf;
pub fn get_base_path() -> PathBuf {
dirs::data_local_dir().unwrap().join("darkfi/app/fud")
}
pub fn get_db_path() -> PathBuf {
dirs::data_local_dir().unwrap().join("darkfi/app/fud/db")
}
pub fn get_downloads_path() -> PathBuf {
dirs::data_local_dir().unwrap().join("darkfi/app/fud/downloads")
}
pub fn get_use_tor_filename() -> PathBuf {
dirs::data_local_dir().unwrap().join("darkfi/app/use_tor.txt")
}
pub fn p2p_datastore_path() -> PathBuf {
dirs::cache_dir().unwrap().join("darkfi/app/fud/p2p")
}
pub fn hostlist_path() -> PathBuf {
dirs::cache_dir().unwrap().join("darkfi/app/fud/hostlist.tsv")
}
}
use paths::*;
macro_rules! t { ($($arg:tt)*) => { trace!(target: "plugin::fud", $($arg)*); } }
macro_rules! d { ($($arg:tt)*) => { debug!(target: "plugin::fud", $($arg)*); } }
macro_rules! i { ($($arg:tt)*) => { info!(target: "plugin::fud", $($arg)*); } }
macro_rules! e { ($($arg:tt)*) => { error!(target: "plugin::fud", $($arg)*); } }
pub type FudPluginPtr = Arc;
pub struct FudPlugin {
node: SceneNodeWeak,
sg_root: SceneNodePtr,
tasks: OnceLock>>,
p2p: P2pPtr,
event_pub: PublisherPtr,
fud: Arc,
tracked_files: Arc>>,
}
impl FudPlugin {
pub async fn new(node: SceneNodeWeak, sg_root: SceneNodePtr, ex: ExecutorPtr) -> Result {
let node_ref = &node.upgrade().unwrap();
// let fud_node_id = PropertyStr::wrap(node_ref, Role::Internal, "node_id", 0).unwrap();
let fud_ready = PropertyBool::wrap(node_ref, Role::Internal, "ready", 0).unwrap();
fud_ready.set(&mut PropertyAtomicGuard::none(), false);
let basedir = get_base_path();
i!("Starting Fud backend");
let db_path = get_db_path();
let db = match Database::open_default(&db_path) {
Ok(db) => db,
Err(err) => {
e!("Kvdb database '{}' failed to open: {err}!", db_path.display());
return Err(Error::KvdbErr)
}
};
let mut fud_settings: FudSettings = Default::default();
fud_settings.base_dir = basedir.to_string_lossy().to_string();
let mut p2p_settings: NetSettings = Default::default();
p2p_settings.magic_bytes = MagicBytes([73, 59, 41, 23]);
p2p_settings.app_version = semver::Version::parse("0.5.0").unwrap();
p2p_settings.app_name = "fud".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://wgxxaifz5gv4iggcflyl67lgmsihffs6bbwobqah4np52t3y3olrnpid.onion:9701",
)
.unwrap(),
);
p2p_settings.seeds.push(
url::Url::parse(
"tor://inx5s3pdzddvgb5ii3oydutmbvw6fvor3oqu65wtxl3pyevtvrdn4had.onion:9701",
)
.unwrap(),
);
p2p_settings.active_profiles = vec!["tor".to_string()];
fud_settings.pow.btc_electrum_nodes.push(
url::Url::parse(
"tor://hezojf7rda2c33yxgcgcvvsxflechdz5vkm64gwlszgx2r4gc5e42kqd.onion:50001",
)
.unwrap(),
);
fud_settings.pow.btc_electrum_nodes.push(
url::Url::parse(
"tor://n4widoxtm3xpo2fjvtdffhb63q5td3utaxkolaegnpzb5khbwxvdrlad.onion:50001",
)
.unwrap(),
);
fud_settings.pow.btc_electrum_nodes.push(
url::Url::parse(
"tor://duras25aqnp3tnn2zgma7pusms6c7umtunyu2sp6e5byotr3c4c6rzad.onion:50001",
)
.unwrap(),
);
fud_settings.pow.btc_electrum_nodes.push(
url::Url::parse(
"tor://n3dz6thzxobyphuosoftgtf36rnsxlsjknke4yrbdys55zvd7nsx7qid.onion:50001",
)
.unwrap(),
);
} 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.active_profiles = vec!["tcp+tls".to_string()];
p2p_settings.seeds.push(url::Url::parse("tcp+tls://lilith0.dark.fi:9700").unwrap());
p2p_settings.seeds.push(url::Url::parse("tcp+tls://lilith1.dark.fi:9700").unwrap());
fud_settings
.pow
.btc_electrum_nodes
.push(url::Url::parse("tcp://fulcrum.grey.pw:50001").unwrap());
fud_settings
.pow
.btc_electrum_nodes
.push(url::Url::parse("tcp://blockstream.info:110").unwrap());
fud_settings
.pow
.btc_electrum_nodes
.push(url::Url::parse("tcp://btc.electroncash.dk:60001").unwrap());
fud_settings
.pow
.btc_electrum_nodes
.push(url::Url::parse("tcp://electrum.direwolfm14.com:50001").unwrap());
fud_settings
.pow
.btc_electrum_nodes
.push(url::Url::parse("tcp://electrum.blockstream.info:50001").unwrap());
}
p2p_settings.p2p_datastore = p2p_datastore_path().into_os_string().into_string().ok();
p2p_settings.hostlist = hostlist_path().into_os_string().into_string().ok();
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)
}
};
p2p.session_direct().start_peer_discovery();
let event_pub = Publisher::new();
let fud: Arc =
match Fud::new(fud_settings, p2p.clone(), &db, event_pub.clone(), ex.clone()).await {
Ok(fud) => fud,
Err(err) => {
e!("Cannot create fud instance: {err}");
return Err(Error::ServiceFailed)
}
};
let self_ = Arc::new(Self {
node: node.clone(),
sg_root,
tasks: OnceLock::new(),
p2p,
event_pub,
fud,
tracked_files: Arc::new(Mutex::new(HashSet::new())),
});
self_.clone().start(ex).await;
Ok(Pimpl::Fud(self_))
}
async fn start(self: Arc, ex: ExecutorPtr) {
i!("Registering Fud protocol");
let registry = self.p2p.protocol_registry();
let fud = self.fud.clone();
let p2p = self.p2p.clone();
registry
.register(SESSION_DIRECT | SESSION_INBOUND, move |channel, _| {
let fud_ = fud.clone();
let p2p_ = p2p.clone();
async move { ProtocolFud::init(fud_, channel, p2p_).await.unwrap() }
})
.await;
let me = Arc::downgrade(&self);
let node = &self.node.upgrade().unwrap();
let method_sub = node.subscribe_method_call("get").unwrap();
let me2 = me.clone();
let get_method_task =
ex.spawn(async move { while Self::process_get(&me2, &method_sub).await {} });
let method_sub = node.subscribe_method_call("track_file").unwrap();
let me2 = me.clone();
let track_file_method_task =
ex.spawn(async move { while Self::process_track_file(&me2, &method_sub).await {} });
let event_pub = self.event_pub.clone();
let me2 = me.clone();
let ev_task = ex.spawn(async move {
Self::process_events(&me2, event_pub).await;
});
let fud = self.fud.clone();
let start_task = ex.spawn(async move {
while fud.start().await.is_err() {
sleep(10).await;
}
});
let tasks = vec![get_method_task, track_file_method_task, ev_task, start_task];
self.tasks.set(tasks).unwrap();
i!("Starting Fud P2P");
while let Err(err) = self.p2p.clone().start().await {
// This usually means we cannot listen on the inbound ports
e!("Failed to start fud's 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;
}
}
fn string_to_hash(str: &str) -> std::io::Result {
let mut hash_buf = vec![];
match bs58::decode(str).onto(&mut hash_buf) {
Ok(_) => {}
Err(_) => {
return Err(std::io::Error::new(std::io::ErrorKind::Other, "Invalid fud hash"))
}
}
if hash_buf.len() != 32 {
return Err(std::io::Error::new(std::io::ErrorKind::Other, "Invalid fud hash"))
}
let mut hash_buf_arr = [0u8; 32];
hash_buf_arr.copy_from_slice(&hash_buf);
Ok(blake3::Hash::from_bytes(hash_buf_arr))
}
fn parse_url(url: &Url) -> std::io::Result<(String, blake3::Hash)> {
let hash_string = url
.host_str()
.map(|s| s.to_string())
.ok_or_else(|| std::io::Error::new(std::io::ErrorKind::Other, "Missing fud hash"))?;
let hash = Self::string_to_hash(&hash_string)?;
Ok((hash_string, hash))
}
fn url_to_file_selection(url: &Url) -> FileSelection {
match url.path() {
"/" | "" => FileSelection::All,
path => {
let mut selection = HashSet::new();
selection.insert(PathBuf::from(path.strip_prefix("/").unwrap_or(path)));
FileSelection::Set(selection)
}
}
}
async fn find_urls_by_hash(&self, hash: &blake3::Hash) -> Vec {
let tracked = self.tracked_files.lock().await;
let hash_str = hash_to_string(hash);
tracked.iter().filter(|url| url.host_str() == Some(hash_str.as_str())).cloned().collect()
}
fn decode_data(
&self,
method_call: &MethodCall,
) -> (Option, std::io::Result<(blake3::Hash, Url, Option)>) {
fn decode_data(data: &[u8]) -> std::io::Result<(String, Url, Option)> {
let mut cur = Cursor::new(&data);
let url = Url::decode(&mut cur)?;
let Some(hash_string) = url.host_str() else {
return Err(std::io::Error::new(std::io::ErrorKind::Other, "Missing fud hash"))
};
let hash_string = hash_string.to_string();
let err_msg = String::decode(&mut cur).ok();
Ok((hash_string, url, err_msg))
}
let Ok((hash_string, url, err_msg)) = decode_data(&method_call.data) else {
return (None, Err(std::io::Error::new(std::io::ErrorKind::Other, "Invalid fud url")))
};
let Ok(hash) = FudPlugin::string_to_hash(&hash_string) else {
return (
Some(hash_string),
Err(std::io::Error::new(std::io::ErrorKind::Other, "Invalid fud url")),
)
};
(Some(hash_string), Ok((hash, url, err_msg)))
}
async fn process_get(me: &Weak, sub: &MethodCallSub) -> bool {
let Ok(method_call) = sub.receive().await else {
d!("Fud event relayer closed");
return false
};
t!("method called: get({method_call:?})");
assert!(method_call.send_res.is_none());
let Some(self_) = me.upgrade() else {
// Should not happen
panic!("self destroyed before get_method_task was stopped!");
};
let (hash_string, data) = self_.decode_data(&method_call);
if let Err(e) = data {
e!("get() method invalid arg data: {e}");
return true
};
let hash_string = hash_string.unwrap();
let (hash, url, _) = data.unwrap();
if self_.node.upgrade().unwrap().get_property_bool("ready").unwrap() {
let file_selection = Self::url_to_file_selection(&url);
let _ = self_
.fud
.get(&hash, &get_downloads_path().join(&hash_string), file_selection)
.await;
}
true
}
/// Get the current file status for a fileurl, a `None` means it should not
/// be updated
async fn get_status(&self, hash: &blake3::Hash, url: &Url) -> Option {
let resources = self.fud.resources().await;
let resource = resources.get(hash);
if resource.is_none() {
return Some(FileMessageStatus::Idle)
}
let resource = resource.unwrap();
let mut path = resource.path.clone();
let file_selection = Self::url_to_file_selection(url);
if let FileSelection::Set(selection) = &file_selection {
if let Some(rel_path) = selection.iter().next() {
path = path.join(rel_path);
}
}
let path = path.to_string_lossy().to_string();
if file_selection.is_disjoint(&resource.last_file_selection) {
return None::
}
let (bytes_downloaded, bytes_total) = self.fud.get_progress(hash, &file_selection).await;
let progress =
if bytes_total != 0 { bytes_downloaded as f32 / bytes_total as f32 * 100. } else { 0. };
match resource.status {
ResourceStatus::Discovering => Some(FileMessageStatus::Downloading { progress }),
ResourceStatus::Downloading => {
if progress < 100. {
Some(FileMessageStatus::Downloading { progress })
} else {
Some(FileMessageStatus::Downloaded { path })
}
}
ResourceStatus::Incomplete(ref err) => {
if progress < 100. {
if let Some(msg) = err {
Some(FileMessageStatus::Error { msg: msg.clone(), progress })
} else {
Some(FileMessageStatus::Error { msg: "incomplete".to_string(), progress })
}
} else {
Some(FileMessageStatus::Downloaded { path })
}
}
ResourceStatus::Verifying => None,
// Seeding status means we have the full resource
// (partial seeding is not supported by fud)
ResourceStatus::Seeding => Some(FileMessageStatus::Downloaded { path }),
}
}
async fn process_track_file(me: &Weak, sub: &MethodCallSub) -> bool {
let Ok(method_call) = sub.receive().await else {
d!("Fud event relayer closed");
return false
};
t!("method called: track_file({method_call:?})");
assert!(method_call.send_res.is_none());
let Some(self_) = me.upgrade() else {
// Should not happen
panic!("self destroyed before track_file_method_task was stopped!");
};
let mut cur = Cursor::new(&method_call.data);
let Ok(url) = Url::decode(&mut cur) else {
e!("track_file() method invalid arg data");
return true
};
self_.track_file(url).await;
true
}
/// Emit file_status_updated signal to all ChatViews
async fn emit_file_status(&self, url: &Url, status: &FileMessageStatus) {
let mut data = vec![];
url.encode(&mut data).unwrap();
status.encode(&mut data).unwrap();
let _ = self.node.upgrade().unwrap().trigger("file_status_updated", data).await;
}
/// Emit error status for a URL
async fn emit_error(&self, url: &Url, msg: String) {
self.emit_file_status(url, &FileMessageStatus::Error { msg, progress: 0. }).await;
}
/// Update tracked files and emit status signal
async fn update_resource(&self, hash: &blake3::Hash) {
let urls = self.find_urls_by_hash(hash).await;
for url in urls {
self.update_fileurl(&url).await;
}
}
async fn update_fileurl(&self, url: &Url) -> bool {
let (_hash_string, hash) = match Self::parse_url(url) {
Ok(h) => h,
Err(err) => {
self.emit_error(url, err.to_string()).await;
return true
}
};
let status = self.get_status(&hash, url).await;
// Emit signal
if let Some(status) = status {
self.emit_file_status(url, &status).await;
return true
}
false
}
/// Emit status for all tracked files
async fn ready_files(&self) {
let tracked = self.tracked_files.lock().await;
let urls: Vec = tracked.iter().cloned().collect();
drop(tracked);
for url in urls {
let (_hash_string, hash) = match Self::parse_url(&url) {
Ok(h) => h,
Err(err) => {
self.emit_error(&url, err.to_string()).await;
continue
}
};
let status = self.get_status(&hash, &url).await;
if let Some(status) = status {
self.emit_file_status(&url, &status).await;
} else {
self.emit_file_status(&url, &FileMessageStatus::Idle).await;
}
}
}
/// Track a file URL (called when the fileurl_detected signal is emitted)
async fn track_file(&self, url: Url) {
let (_hash_string, _hash) = match Self::parse_url(&url) {
Ok(h) => h,
Err(err) => {
self.emit_error(&url, err.to_string()).await;
return
}
};
if self.node.upgrade().unwrap().get_property_bool("ready").unwrap() {
let updated = self.update_fileurl(&url).await;
if !updated {
self.emit_file_status(&url, &FileMessageStatus::Idle).await;
}
}
let mut tracked = self.tracked_files.lock().await;
tracked.insert(url);
}
async fn process_events(me: &Weak, publisher: PublisherPtr) {
let Some(self_) = me.upgrade() else {
// Should not happen
panic!("self destroyed before ev_task was stopped!");
};
let sub = publisher.subscribe().await;
loop {
match sub.receive().await {
FudEvent::Ready => {
let atom = &mut PropertyAtomicGuard::none();
self_
.node
.upgrade()
.unwrap()
.set_property_bool(atom, Role::App, "ready", true)
.unwrap();
self_.ready_files().await;
}
FudEvent::DownloadStarted(ev) => {
self_.update_resource(&ev.resource.hash).await;
}
FudEvent::ChunkDownloadCompleted(ev) => {
self_.update_resource(&ev.resource.hash).await;
}
FudEvent::DownloadCompleted(ev) => {
self_.update_resource(&ev.resource.hash).await;
}
FudEvent::ResourceUpdated(ev) => {
self_.update_resource(&ev.resource.hash).await;
}
FudEvent::DownloadError(ev) => {
self_.update_resource(&ev.hash).await;
}
FudEvent::MissingChunks(ev) => {
self_.update_resource(&ev.hash).await;
}
FudEvent::MetadataNotFound(ev) => {
self_.update_resource(&ev.hash).await;
}
_ => {}
};
}
}
}