jsonrpc.rs 28 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2025 Dyne.org foundation
  4. *
  5. * This program is free software: you can redistribute it and/or modify
  6. * it under the terms of the GNU Affero General Public License as
  7. * published by the Free Software Foundation, either version 3 of the
  8. * License, or (at your option) any later version.
  9. *
  10. * This program is distributed in the hope that it will be useful,
  11. * but WITHOUT ANY WARRANTY; without even the implied warranty of
  12. * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
  13. * GNU Affero General Public License for more details.
  14. *
  15. * You should have received a copy of the GNU Affero General Public License
  16. * along with this program. If not, see <https://www.gnu.org/licenses/>.
  17. */
  18. use std::{
  19. collections::{HashMap, HashSet},
  20. fs::create_dir_all,
  21. path::PathBuf,
  22. sync::Arc,
  23. };
  24. use async_trait::async_trait;
  25. use smol::lock::{Mutex, MutexGuard};
  26. use tinyjson::JsonValue;
  27. use tracing::{debug, info, warn};
  28. use darkfi::{
  29. event_graph::EventGraphPtr,
  30. net,
  31. rpc::{
  32. jsonrpc::{ErrorCode, JsonError, JsonRequest, JsonResult, JsonSubscriber},
  33. p2p_method::HandlerP2p,
  34. server::RequestHandler,
  35. },
  36. system::StoppableTaskPtr,
  37. util::{path::expand_path, time::Timestamp},
  38. Error,
  39. };
  40. use taud::{
  41. error::{to_json_result, TaudError, TaudResult},
  42. month_tasks::MonthTasks,
  43. task_info::{Comment, TaskInfo},
  44. util::set_event,
  45. };
  46. use crate::Workspace;
  47. const DEFAULT_WORKSPACE: &str = "darkfi-dev";
  48. pub struct JsonRpcInterface {
  49. dataset_path: PathBuf,
  50. notify_queue_sender: smol::channel::Sender<TaskInfo>,
  51. nickname: String,
  52. workspace: Mutex<String>,
  53. workspaces: Arc<HashMap<String, Workspace>>,
  54. p2p: net::P2pPtr,
  55. event_graph: EventGraphPtr,
  56. dnet_sub: JsonSubscriber,
  57. deg_sub: JsonSubscriber,
  58. rpc_connections: Mutex<HashSet<StoppableTaskPtr>>,
  59. }
  60. #[async_trait]
  61. impl RequestHandler<()> for JsonRpcInterface {
  62. async fn handle_request(&self, req: JsonRequest) -> JsonResult {
  63. let rep = match req.method.as_str() {
  64. "add" => self.add(req.params).await,
  65. "get_ref_ids" => self.get_ref_ids(req.params).await,
  66. "get_archive_ref_ids" => self.get_archive_ref_ids(req.params).await,
  67. "modify" => self.modify(req.params).await,
  68. "set_state" => self.set_state(req.params).await,
  69. "set_comment" => self.set_comment(req.params).await,
  70. "get_task_by_ref_id" => self.get_task_by_ref_id(req.params).await,
  71. "switch_ws" => self.switch_ws(req.params).await,
  72. "get_ws" => self.get_ws(req.params).await,
  73. "export" => self.export_to(req.params).await,
  74. "import" => self.import_from(req.params).await,
  75. "fetch_deactive_tasks" => self.fetch_deactive_tasks(req.params).await,
  76. "fetch_archive_task" => self.fetch_archive_task(req.params).await,
  77. "ping" => return self.pong(req.id, req.params).await,
  78. "dnet.subscribe_events" => return self.dnet_subscribe_events(req.id, req.params).await,
  79. "dnet.switch" => self.dnet_switch(req.params).await,
  80. "deg.switch" => self.deg_switch(req.id, req.params).await,
  81. "deg.subscribe_events" => return self.deg_subscribe_events(req.id, req.params).await,
  82. "eventgraph.get_info" => return self.eg_get_info(req.id, req.params).await,
  83. "p2p.get_info" => return self.p2p_get_info(req.id, req.params).await,
  84. _ => return JsonError::new(ErrorCode::MethodNotFound, None, req.id).into(),
  85. };
  86. to_json_result(rep, req.id)
  87. }
  88. async fn connections_mut(&self) -> MutexGuard<'life0, HashSet<StoppableTaskPtr>> {
  89. self.rpc_connections.lock().await
  90. }
  91. }
  92. impl HandlerP2p for JsonRpcInterface {
  93. fn p2p(&self) -> net::P2pPtr {
  94. self.p2p.clone()
  95. }
  96. }
  97. impl JsonRpcInterface {
  98. #[allow(clippy::too_many_arguments)]
  99. pub fn new(
  100. dataset_path: PathBuf,
  101. notify_queue_sender: smol::channel::Sender<TaskInfo>,
  102. nickname: String,
  103. workspaces: Arc<HashMap<String, Workspace>>,
  104. p2p: net::P2pPtr,
  105. event_graph: EventGraphPtr,
  106. dnet_sub: JsonSubscriber,
  107. deg_sub: JsonSubscriber,
  108. ) -> Self {
  109. let workspace = Mutex::new(DEFAULT_WORKSPACE.to_string());
  110. Self {
  111. dataset_path,
  112. nickname,
  113. workspace,
  114. workspaces,
  115. notify_queue_sender,
  116. p2p,
  117. event_graph,
  118. rpc_connections: Mutex::new(HashSet::new()),
  119. dnet_sub,
  120. deg_sub,
  121. }
  122. }
  123. // RPCAPI:
  124. // Activate or deactivate dnet in the P2P stack.
  125. // By sending `true`, dnet will be activated, and by sending `false` dnet will
  126. // be deactivated. Returns `true` on success.
  127. //
  128. // --> {"jsonrpc": "2.0", "method": "dnet_switch", "params": [true], "id": 42}
  129. // <-- {"jsonrpc": "2.0", "result": true, "id": 42}
  130. async fn dnet_switch(&self, params: JsonValue) -> TaudResult<JsonValue> {
  131. let params = params.get::<Vec<JsonValue>>().unwrap();
  132. if params.len() != 1 || !params[0].is_bool() {
  133. return Err(TaudError::InvalidData("Invalid parameters".into()))
  134. }
  135. let switch = params[0].get::<bool>().unwrap();
  136. if *switch {
  137. self.p2p.dnet_enable();
  138. } else {
  139. self.p2p.dnet_disable();
  140. }
  141. Ok(JsonValue::Boolean(true))
  142. }
  143. // RPCAPI:
  144. // Initializes a subscription to p2p dnet events.
  145. // Once a subscription is established, `darkirc` will send JSON-RPC notifications of
  146. // new network events to the subscriber.
  147. //
  148. // --> {"jsonrpc": "2.0", "method": "dnet.subscribe_events", "params": [], "id": 1}
  149. // <-- {"jsonrpc": "2.0", "method": "dnet.subscribe_events", "params": [`event`]}
  150. pub async fn dnet_subscribe_events(&self, id: u16, params: JsonValue) -> JsonResult {
  151. let params = params.get::<Vec<JsonValue>>().unwrap();
  152. if !params.is_empty() {
  153. return JsonError::new(ErrorCode::InvalidParams, None, id).into()
  154. }
  155. self.dnet_sub.clone().into()
  156. }
  157. // RPCAPI:
  158. // Initializes a subscription to deg events.
  159. // Once a subscription is established, apps using eventgraph will send JSON-RPC notifications of
  160. // new eventgraph events to the subscriber.
  161. //
  162. // --> {"jsonrpc": "2.0", "method": "deg.subscribe_events", "params": [], "id": 1}
  163. // <-- {"jsonrpc": "2.0", "method": "deg.subscribe_events", "params": [`event`]}
  164. pub async fn deg_subscribe_events(&self, id: u16, params: JsonValue) -> JsonResult {
  165. let params = params.get::<Vec<JsonValue>>().unwrap();
  166. if !params.is_empty() {
  167. return JsonError::new(ErrorCode::InvalidParams, None, id).into()
  168. }
  169. self.deg_sub.clone().into()
  170. }
  171. // RPCAPI:
  172. // Activate or deactivate deg in the EVENTGRAPH.
  173. // By sending `true`, deg will be activated, and by sending `false` deg
  174. // will be deactivated. Returns `true` on success.
  175. //
  176. // --> {"jsonrpc": "2.0", "method": "deg.switch", "params": [true], "id": 42}
  177. // <-- {"jsonrpc": "2.0", "result": true, "id": 42}
  178. async fn deg_switch(&self, _id: u16, params: JsonValue) -> TaudResult<JsonValue> {
  179. let params = params.get::<Vec<JsonValue>>().unwrap();
  180. if params.len() != 1 || !params[0].is_bool() {
  181. return Err(TaudError::InvalidData("Invalid parameters".into()))
  182. }
  183. let switch = params[0].get::<bool>().unwrap();
  184. if *switch {
  185. self.event_graph.deg_enable().await;
  186. } else {
  187. self.event_graph.deg_disable().await;
  188. }
  189. Ok(JsonValue::Boolean(true))
  190. }
  191. // RPCAPI:
  192. // Get EVENTGRAPH info.
  193. //
  194. // --> {"jsonrpc": "2.0", "method": "deg.switch", "params": [true], "id": 42}
  195. // <-- {"jsonrpc": "2.0", "result": true, "id": 42}
  196. async fn eg_get_info(&self, id: u16, params: JsonValue) -> JsonResult {
  197. let params_ = params.get::<Vec<JsonValue>>().unwrap();
  198. if !params_.is_empty() {
  199. return JsonError::new(ErrorCode::InvalidParams, None, id).into()
  200. }
  201. self.event_graph.eventgraph_info(id, params).await
  202. }
  203. // RPCAPI:
  204. // Add new task and returns `true` upon success.
  205. // --> {"jsonrpc": "2.0", "method": "add",
  206. // "params":
  207. // [{
  208. // "title": "..",
  209. // "desc": "..",
  210. // assign: [..],
  211. // project: [..],
  212. // "due": ..,
  213. // "rank": ..
  214. // }],
  215. // "id": 1
  216. // }
  217. // <-- {"jsonrpc": "2.0", "result": true, "id": 1}
  218. async fn add(&self, params: JsonValue) -> TaudResult<JsonValue> {
  219. let params = params.get::<Vec<JsonValue>>().unwrap();
  220. debug!(target: "tau", "JsonRpc::add() params {params:?}");
  221. if !params[0].is_object() {
  222. return Err(TaudError::InvalidData("Invalid parameters".to_string()))
  223. }
  224. let params = params[0].get::<HashMap<String, JsonValue>>().unwrap();
  225. if params.len() != 9 {
  226. return Err(TaudError::InvalidData("Invalid parameters".to_string()))
  227. }
  228. let due = match params["due"] {
  229. JsonValue::Null => None,
  230. JsonValue::Number(numba) => Some(Timestamp::from_u64(numba as u64)),
  231. _ => return Err(TaudError::InvalidData("Invalid parameter \"due\"".to_string())),
  232. };
  233. let rank = match params["rank"] {
  234. JsonValue::Null => None,
  235. JsonValue::Number(numba) => Some(numba as f32),
  236. _ => return Err(TaudError::InvalidData("Invalid parameter \"rank\"".to_string())),
  237. };
  238. let tags = {
  239. let mut tags = vec![];
  240. for val in params["tags"].get::<Vec<JsonValue>>().unwrap().iter() {
  241. if let Some(tag) = val.get::<String>() {
  242. tags.push(tag.clone());
  243. } else {
  244. return Err(TaudError::InvalidData("Invalid parameter \"tags\"".to_string()))
  245. }
  246. }
  247. tags
  248. };
  249. let assigns = {
  250. let mut assigns = vec![];
  251. for val in params["assign"].get::<Vec<JsonValue>>().unwrap().iter() {
  252. if let Some(assign) = val.get::<String>() {
  253. assigns.push(assign.clone());
  254. } else {
  255. return Err(TaudError::InvalidData("Invalid parameter \"assign\"".to_string()))
  256. }
  257. }
  258. assigns
  259. };
  260. let projects = {
  261. let mut projects = vec![];
  262. for val in params["project"].get::<Vec<JsonValue>>().unwrap().iter() {
  263. if let Some(project) = val.get::<String>() {
  264. projects.push(project.clone());
  265. } else {
  266. return Err(TaudError::InvalidData("Invalid parameter \"project\"".to_string()))
  267. }
  268. }
  269. projects
  270. };
  271. let created_at = match params["created_at"] {
  272. JsonValue::Number(numba) => Some(numba as u64),
  273. _ => return Err(TaudError::InvalidData("Invalid parameter \"created_at\"".to_string())),
  274. };
  275. let ws = self.workspace.lock().await.clone();
  276. if self.workspaces.get(&ws).unwrap().write_key.is_none() {
  277. info!("You don't have write access!");
  278. return Ok(JsonValue::Boolean(false))
  279. }
  280. let mut new_task: TaskInfo = TaskInfo::new(
  281. ws,
  282. params["title"].get::<String>().unwrap(),
  283. params["desc"].get::<String>().unwrap(),
  284. &self.nickname,
  285. due,
  286. rank,
  287. Timestamp::from_u64(created_at.unwrap()),
  288. )?;
  289. new_task.set_project(&projects);
  290. new_task.set_assign(&assigns);
  291. new_task.set_tags(&tags);
  292. self.notify_queue_sender.send(new_task.clone()).await.map_err(Error::from)?;
  293. Ok(new_task.ref_id.clone().into())
  294. }
  295. // RPCAPI:
  296. // List tasks
  297. // --> {"jsonrpc": "2.0", "method": "get_ids", "params": [], "id": 1}
  298. // <-- {"jsonrpc": "2.0", "result": [task_id, ...], "id": 1}
  299. async fn get_ref_ids(&self, params: JsonValue) -> TaudResult<JsonValue> {
  300. let params = params.get::<Vec<JsonValue>>().unwrap();
  301. debug!(target: "tau", "JsonRpc::get_ids() params {params:?}");
  302. let ws = self.workspace.lock().await.clone();
  303. let tasks = MonthTasks::load_current_tasks(&self.dataset_path, ws, false)?;
  304. let task_ref_ids: Vec<JsonValue> =
  305. tasks.iter().map(|task| JsonValue::String(task.get_ref_id())).collect();
  306. Ok(JsonValue::Array(task_ref_ids))
  307. }
  308. // RPCAPI:
  309. // List tasks
  310. // --> {"jsonrpc": "2.0", "method": "get_ids", "params": [], "id": 1}
  311. // <-- {"jsonrpc": "2.0", "result": [task_id, ...], "id": 1}
  312. async fn get_archive_ref_ids(&self, params: JsonValue) -> TaudResult<JsonValue> {
  313. let params = params.get::<Vec<JsonValue>>().unwrap();
  314. debug!(target: "tau", "JsonRpc::get_archive_ref_ids() params {params:?}");
  315. let month = match params[0].get::<String>() {
  316. Some(u64_str) => match u64_str.parse::<u64>() {
  317. Ok(v) => Some(Timestamp::from_u64(v)),
  318. //Err(e) => return Err(TaudError::InvalidData(e.to_string())),
  319. Err(_) => None,
  320. },
  321. None => None,
  322. };
  323. let ws = self.workspace.lock().await.clone();
  324. let tasks = MonthTasks::load_stop_tasks(&self.dataset_path, ws, month.as_ref())?;
  325. let task_ref_ids: Vec<JsonValue> =
  326. tasks.iter().map(|task| JsonValue::String(task.get_ref_id())).collect();
  327. Ok(JsonValue::Array(task_ref_ids))
  328. }
  329. // RPCAPI:
  330. // Modify task and returns `true` upon success.
  331. // --> {"jsonrpc": "2.0", "method": "modify", "params": [task_id, {"title": "new title"} ], "id": 1}
  332. // <-- {"jsonrpc": "2.0", "result": true, "id": 1}
  333. async fn modify(&self, params: JsonValue) -> TaudResult<JsonValue> {
  334. let params = params.get::<Vec<JsonValue>>().unwrap();
  335. debug!(target: "tau", "JsonRpc::modify() params {params:?}");
  336. if params.len() != 2 || !params[0].is_string() || !params[1].is_object() {
  337. return Err(TaudError::InvalidData("len of params should be 2".into()))
  338. }
  339. let ws = self.workspace.lock().await.clone();
  340. if self.workspaces.get(&ws).unwrap().write_key.is_none() {
  341. info!("You don't have write access!");
  342. return Ok(JsonValue::Boolean(false))
  343. }
  344. let task = self.check_params_for_modify(
  345. params[0].get::<String>().unwrap(),
  346. params[1].get::<HashMap<String, JsonValue>>().unwrap(),
  347. ws,
  348. )?;
  349. self.notify_queue_sender.send(task).await.map_err(Error::from)?;
  350. Ok(JsonValue::Boolean(true))
  351. }
  352. // RPCAPI:
  353. // Set state for a task and returns `true` upon success.
  354. // --> {"jsonrpc": "2.0", "method": "set_state", "params": [task_id, state], "id": 1}
  355. // <-- {"jsonrpc": "2.0", "result": true, "id": 1}
  356. async fn set_state(&self, params: JsonValue) -> TaudResult<JsonValue> {
  357. // Allowed states for a task
  358. let states = ["stop", "start", "open", "pause"];
  359. let params = params.get::<Vec<JsonValue>>().unwrap();
  360. debug!(target: "tau", "JsonRpc::set_state() params {params:?}");
  361. if params.len() != 2 || !params[0].is_string() || !params[1].is_string() {
  362. return Err(TaudError::InvalidData("len of params should be 2".into()))
  363. }
  364. let state = params[1].get::<String>().unwrap();
  365. let ws = self.workspace.lock().await.clone();
  366. if self.workspaces.get(&ws).unwrap().write_key.is_none() {
  367. info!("You don't have write access!");
  368. return Ok(JsonValue::Boolean(false))
  369. }
  370. let mut task: TaskInfo =
  371. self.load_task_by_ref_id(params[0].get::<String>().unwrap(), ws)?;
  372. if states.contains(&state.as_str()) {
  373. task.set_state(state);
  374. set_event(&mut task, "state", &self.nickname, state);
  375. }
  376. self.notify_queue_sender.send(task).await.map_err(Error::from)?;
  377. Ok(JsonValue::Boolean(true))
  378. }
  379. // RPCAPI:
  380. // Set comment for a task and returns `true` upon success.
  381. // --> {"jsonrpc": "2.0", "method": "set_comment", "params": [task_id, comment_content], "id": 1}
  382. // <-- {"jsonrpc": "2.0", "result": true, "id": 1}
  383. async fn set_comment(&self, params: JsonValue) -> TaudResult<JsonValue> {
  384. let params = params.get::<Vec<JsonValue>>().unwrap();
  385. debug!(target: "tau", "JsonRpc::set_comment() params {params:?}");
  386. if params.len() != 2 || !params[0].is_string() || !params[1].is_string() {
  387. return Err(TaudError::InvalidData("len of params should be 2".into()))
  388. }
  389. let ref_id = params[0].get::<String>().unwrap();
  390. let comment_content = params[1].get::<String>().unwrap();
  391. let ws = self.workspace.lock().await.clone();
  392. if self.workspaces.get(&ws).unwrap().write_key.is_none() {
  393. info!("You don't have write access!");
  394. return Ok(JsonValue::Boolean(false))
  395. }
  396. let mut task: TaskInfo = self.load_task_by_ref_id(ref_id, ws)?;
  397. task.set_comment(Comment::new(comment_content, &self.nickname));
  398. set_event(&mut task, "comment", &self.nickname, comment_content);
  399. self.notify_queue_sender.send(task).await.map_err(Error::from)?;
  400. Ok(JsonValue::Boolean(true))
  401. }
  402. // RPCAPI:
  403. // Get a task by id.
  404. // --> {"jsonrpc": "2.0", "method": "get_task_by_id", "params": [task_id], "id": 1}
  405. // <-- {"jsonrpc": "2.0", "result": "task", "id": 1}
  406. async fn get_task_by_ref_id(&self, params: JsonValue) -> TaudResult<JsonValue> {
  407. let params = params.get::<Vec<JsonValue>>().unwrap();
  408. debug!(target: "tau", "JsonRpc::get_task_by_ref_id() params {params:?}");
  409. if params.len() != 1 || !params[0].is_string() {
  410. return Err(TaudError::InvalidData("len of params should be 1".into()))
  411. }
  412. let ws = self.workspace.lock().await.clone();
  413. let task: TaskInfo = self.load_task_by_ref_id(params[0].get::<String>().unwrap(), ws)?;
  414. let task: JsonValue = (&task).into();
  415. Ok(task)
  416. }
  417. // RPCAPI:
  418. // Get all tasks.
  419. // --> {"jsonrpc": "2.0", "method": "fetch_deactive_tasks", "params": [task_id], "id": 1}
  420. // <-- {"jsonrpc": "2.0", "result": "task", "id": 1}
  421. async fn fetch_deactive_tasks(&self, params: JsonValue) -> TaudResult<JsonValue> {
  422. let params = params.get::<Vec<JsonValue>>().unwrap();
  423. debug!(target: "tau", "JsonRpc::fetch_deactive_tasks() params {params:?}");
  424. if params.len() != 1 || !params[0].is_string() {
  425. return Err(TaudError::InvalidData("len of params should be 1".into()))
  426. }
  427. let month = match params[0].get::<String>() {
  428. Some(u64_str) => match u64_str.parse::<u64>() {
  429. Ok(v) => Some(Timestamp::from_u64(v)),
  430. //Err(e) => return Err(TaudError::InvalidData(e.to_string())),
  431. Err(_) => None,
  432. },
  433. None => None,
  434. };
  435. let ws = self.workspace.lock().await.clone();
  436. let tasks = MonthTasks::load_stop_tasks(&self.dataset_path, ws, month.as_ref())?;
  437. let tasks: Vec<JsonValue> = tasks.iter().map(|x| x.into()).collect();
  438. Ok(JsonValue::Array(tasks))
  439. }
  440. async fn fetch_archive_task(&self, params: JsonValue) -> TaudResult<JsonValue> {
  441. let params = params.get::<Vec<JsonValue>>().unwrap();
  442. debug!(target: "tau", "JsonRpc::fetch_archive_task() params {params:?}");
  443. if params.len() != 2 || !params[0].is_string() || !params[1].is_string() {
  444. return Err(TaudError::InvalidData("len of params should be 2".into()))
  445. }
  446. let ref_id = params[0].get::<String>().unwrap();
  447. let month = match params[1].get::<String>() {
  448. Some(u64_str) => match u64_str.parse::<u64>() {
  449. Ok(v) => Some(Timestamp::from_u64(v)),
  450. //Err(e) => return Err(TaudError::InvalidData(e.to_string())),
  451. Err(_) => None,
  452. },
  453. None => None,
  454. };
  455. let ws = self.workspace.lock().await.clone();
  456. let mut tasks = MonthTasks::load_stop_tasks(&self.dataset_path, ws, month.as_ref())?;
  457. tasks.retain(|x| x.ref_id == *ref_id);
  458. if tasks.len() != 1 {
  459. return Err(TaudError::InvalidData("Must return a single value".into()))
  460. }
  461. let task: JsonValue = (&tasks[0]).into();
  462. Ok(task)
  463. }
  464. // RPCAPI:
  465. // Switch tasks workspace.
  466. // --> {"jsonrpc": "2.0", "method": "switch_ws", "params": [workspace], "id": 1}
  467. // <-- {"jsonrpc": "2.0", "result": "true", "id": 1}
  468. async fn switch_ws(&self, params: JsonValue) -> TaudResult<JsonValue> {
  469. let params = params.get::<Vec<JsonValue>>().unwrap();
  470. debug!(target: "tau", "JsonRpc::switch_ws() params {params:?}");
  471. if params.len() != 1 {
  472. return Err(TaudError::InvalidData("len of params should be 1".into()))
  473. }
  474. if !params[0].is_string() {
  475. return Err(TaudError::InvalidData("Invalid workspace".into()))
  476. }
  477. let ws = params[0].get::<String>().unwrap();
  478. let mut s = self.workspace.lock().await;
  479. if self.workspaces.contains_key(ws) {
  480. *s = ws.to_string()
  481. } else {
  482. warn!("Workspace \"{ws}\" is not configured");
  483. return Ok(JsonValue::Boolean(false))
  484. }
  485. Ok(JsonValue::Boolean(true))
  486. }
  487. // RPCAPI:
  488. // Get workspace.
  489. // --> {"jsonrpc": "2.0", "method": "get_ws", "params": [], "id": 1}
  490. // <-- {"jsonrpc": "2.0", "result": "workspace", "id": 1}
  491. async fn get_ws(&self, params: JsonValue) -> TaudResult<JsonValue> {
  492. let params = params.get::<Vec<JsonValue>>().unwrap();
  493. debug!(target: "tau", "JsonRpc::get_ws() params {params:?}");
  494. let ws = self.workspace.lock().await.clone();
  495. Ok(JsonValue::String(ws))
  496. }
  497. // RPCAPI:
  498. // Export tasks.
  499. // --> {"jsonrpc": "2.0", "method": "export_to", "params": [path], "id": 1}
  500. // <-- {"jsonrpc": "2.0", "result": "true", "id": 1}
  501. async fn export_to(&self, params: JsonValue) -> TaudResult<JsonValue> {
  502. let params = params.get::<Vec<JsonValue>>().unwrap();
  503. debug!(target: "tau", "JsonRpc::export_to() params {params:?}");
  504. if params.len() != 1 {
  505. return Err(TaudError::InvalidData("len of params should be 1".into()))
  506. }
  507. if !params[0].is_string() {
  508. return Err(TaudError::InvalidData("Invalid path".into()))
  509. }
  510. // mkdir datastore_path if not exists
  511. let path = params[0].get::<String>().unwrap();
  512. let path = expand_path(path)?.join("exported_tasks");
  513. create_dir_all(path.join("month")).map_err(Error::from)?;
  514. create_dir_all(path.join("task")).map_err(Error::from)?;
  515. let ws = self.workspace.lock().await.clone();
  516. let tasks = MonthTasks::load_current_tasks(&self.dataset_path, ws, true)?;
  517. for task in tasks {
  518. task.save(&path)?;
  519. }
  520. Ok(JsonValue::Boolean(true))
  521. }
  522. // RPCAPI:
  523. // Import tasks.
  524. // --> {"jsonrpc": "2.0", "method": "import_from", "params": [path], "id": 1}
  525. // <-- {"jsonrpc": "2.0", "result": "true", "id": 1}
  526. async fn import_from(&self, params: JsonValue) -> TaudResult<JsonValue> {
  527. let params = params.get::<Vec<JsonValue>>().unwrap();
  528. debug!(target: "tau", "JsonRpc::import_from() params {params:?}");
  529. if params.len() != 1 {
  530. return Err(TaudError::InvalidData("len of params should be 1".into()))
  531. }
  532. if !params[0].is_string() {
  533. return Err(TaudError::InvalidData("Invalid path".into()))
  534. }
  535. let path = params[0].get::<String>().unwrap();
  536. let path = expand_path(path)?.join("exported_tasks");
  537. let ws = self.workspace.lock().await.clone();
  538. if self.workspaces.get(&ws).unwrap().write_key.is_none() {
  539. info!("You don't have write access!");
  540. return Ok(JsonValue::Boolean(false))
  541. }
  542. let imported_tasks = MonthTasks::load_current_tasks(&path, ws.clone(), true)?;
  543. for task in imported_tasks {
  544. if MonthTasks::load_current_tasks(&self.dataset_path, ws.clone(), false)?
  545. .into_iter()
  546. .map(|t| t.ref_id)
  547. .any(|x| x == task.ref_id)
  548. {
  549. continue
  550. }
  551. self.notify_queue_sender.send(task).await.map_err(Error::from)?;
  552. }
  553. Ok(JsonValue::Boolean(true))
  554. }
  555. fn load_task_by_ref_id(&self, task_ref_id: &str, ws: String) -> TaudResult<TaskInfo> {
  556. let tasks = MonthTasks::load_current_tasks(&self.dataset_path, ws, false)?;
  557. let task = tasks.into_iter().find(|t| (t.get_ref_id()) == task_ref_id);
  558. task.ok_or(TaudError::InvalidId)
  559. }
  560. fn check_params_for_modify(
  561. &self,
  562. task_ref_id: &str,
  563. fields: &HashMap<String, JsonValue>,
  564. ws: String,
  565. ) -> TaudResult<TaskInfo> {
  566. let mut task: TaskInfo = self.load_task_by_ref_id(task_ref_id, ws)?;
  567. if fields.contains_key("title") {
  568. let title = fields["title"].get::<String>().unwrap();
  569. if !title.is_empty() {
  570. task.set_title(title);
  571. set_event(&mut task, "title", &self.nickname, title);
  572. }
  573. }
  574. if fields.contains_key("desc") {
  575. let desc = fields["desc"].get::<String>().unwrap();
  576. if !desc.is_empty() {
  577. task.set_desc(desc);
  578. set_event(&mut task, "desc", &self.nickname, desc);
  579. }
  580. }
  581. if fields.contains_key("rank") {
  582. match fields["rank"] {
  583. JsonValue::Null => set_event(&mut task, "rank", &self.nickname, "None"),
  584. JsonValue::Number(rank) => {
  585. task.set_rank(Some(rank as f32));
  586. set_event(&mut task, "rank", &self.nickname, &rank.to_string())
  587. }
  588. _ => unreachable!(),
  589. }
  590. }
  591. if fields.contains_key("due") {
  592. match &fields["due"] {
  593. JsonValue::Null => set_event(&mut task, "due", &self.nickname, "None"),
  594. JsonValue::Number(ts_num) => {
  595. task.set_due(Some(Timestamp::from_u64(*ts_num as u64)));
  596. set_event(&mut task, "due", &self.nickname, &ts_num.to_string())
  597. }
  598. _ => unreachable!(),
  599. }
  600. }
  601. if fields.contains_key("assign") {
  602. let assign: Vec<String> = fields["assign"]
  603. .get::<Vec<JsonValue>>()
  604. .unwrap()
  605. .iter()
  606. .map(|x| x.get::<String>().unwrap().clone())
  607. .collect();
  608. if !assign.is_empty() {
  609. task.set_assign(&assign);
  610. set_event(&mut task, "assign", &self.nickname, &assign.join(", "));
  611. }
  612. }
  613. if fields.contains_key("project") {
  614. let project: Vec<String> = fields["project"]
  615. .get::<Vec<JsonValue>>()
  616. .unwrap()
  617. .iter()
  618. .map(|x| x.get::<String>().unwrap().clone())
  619. .collect();
  620. if !project.is_empty() {
  621. task.set_project(&project);
  622. set_event(&mut task, "project", &self.nickname, &project.join(", "));
  623. }
  624. }
  625. if fields.contains_key("tags") {
  626. let tags: Vec<String> = fields["tags"]
  627. .get::<Vec<JsonValue>>()
  628. .unwrap()
  629. .iter()
  630. .map(|x| x.get::<String>().unwrap().clone())
  631. .collect();
  632. if !tags.is_empty() {
  633. task.set_tags(&tags);
  634. set_event(&mut task, "tags", &self.nickname, &tags.join(", "));
  635. }
  636. }
  637. Ok(task)
  638. }
  639. }