jsonrpc.rs 28 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2024 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 crypto_box::ChaChaBox;
  26. use log::{debug, warn};
  27. use smol::lock::{Mutex, MutexGuard};
  28. use tinyjson::JsonValue;
  29. use darkfi::{
  30. event_graph::EventGraphPtr,
  31. net,
  32. rpc::{
  33. jsonrpc::{ErrorCode, JsonError, JsonRequest, JsonResult, JsonSubscriber},
  34. p2p_method::HandlerP2p,
  35. server::RequestHandler,
  36. },
  37. system::StoppableTaskPtr,
  38. util::{path::expand_path, time::Timestamp},
  39. Error,
  40. };
  41. use taud::{
  42. error::{to_json_result, TaudError, TaudResult},
  43. month_tasks::MonthTasks,
  44. task_info::{Comment, TaskInfo},
  45. util::{check_write_access, set_event},
  46. };
  47. pub struct JsonRpcInterface {
  48. dataset_path: PathBuf,
  49. notify_queue_sender: smol::channel::Sender<TaskInfo>,
  50. nickname: String,
  51. workspace: Mutex<String>,
  52. workspaces: Arc<HashMap<String, ChaChaBox>>,
  53. write: Option<String>,
  54. password: Option<String>,
  55. p2p: net::P2pPtr,
  56. event_graph: EventGraphPtr,
  57. dnet_sub: JsonSubscriber,
  58. deg_sub: JsonSubscriber,
  59. rpc_connections: Mutex<HashSet<StoppableTaskPtr>>,
  60. }
  61. #[async_trait]
  62. impl RequestHandler for JsonRpcInterface {
  63. async fn handle_request(&self, req: JsonRequest) -> JsonResult {
  64. let rep = match req.method.as_str() {
  65. "add" => self.add(req.params).await,
  66. "get_ref_ids" => self.get_ref_ids(req.params).await,
  67. "get_archive_ref_ids" => self.get_archive_ref_ids(req.params).await,
  68. "modify" => self.modify(req.params).await,
  69. "set_state" => self.set_state(req.params).await,
  70. "set_comment" => self.set_comment(req.params).await,
  71. "get_task_by_ref_id" => self.get_task_by_ref_id(req.params).await,
  72. "switch_ws" => self.switch_ws(req.params).await,
  73. "get_ws" => self.get_ws(req.params).await,
  74. "export" => self.export_to(req.params).await,
  75. "import" => self.import_from(req.params).await,
  76. "fetch_deactive_tasks" => self.fetch_deactive_tasks(req.params).await,
  77. "fetch_archive_task" => self.fetch_archive_task(req.params).await,
  78. "ping" => return self.pong(req.id, req.params).await,
  79. "dnet.subscribe_events" => return self.dnet_subscribe_events(req.id, req.params).await,
  80. "dnet.switch" => self.dnet_switch(req.params).await,
  81. "deg.switch" => self.deg_switch(req.id, req.params).await,
  82. "deg.subscribe_events" => return self.deg_subscribe_events(req.id, req.params).await,
  83. "eventgraph.get_info" => return self.eg_get_info(req.id, req.params).await,
  84. // TODO: make this optional
  85. "p2p.get_info" => return self.p2p_get_info(req.id, req.params).await,
  86. _ => return JsonError::new(ErrorCode::MethodNotFound, None, req.id).into(),
  87. };
  88. to_json_result(rep, req.id)
  89. }
  90. async fn connections_mut(&self) -> MutexGuard<'_, HashSet<StoppableTaskPtr>> {
  91. self.rpc_connections.lock().await
  92. }
  93. }
  94. impl HandlerP2p for JsonRpcInterface {
  95. fn p2p(&self) -> net::P2pPtr {
  96. self.p2p.clone()
  97. }
  98. }
  99. impl JsonRpcInterface {
  100. #[allow(clippy::too_many_arguments)]
  101. pub fn new(
  102. dataset_path: PathBuf,
  103. notify_queue_sender: smol::channel::Sender<TaskInfo>,
  104. nickname: String,
  105. workspaces: Arc<HashMap<String, ChaChaBox>>,
  106. write: Option<String>,
  107. password: Option<String>,
  108. p2p: net::P2pPtr,
  109. event_graph: EventGraphPtr,
  110. dnet_sub: JsonSubscriber,
  111. deg_sub: JsonSubscriber,
  112. ) -> Self {
  113. let workspace = Mutex::new(workspaces.iter().last().unwrap().0.clone());
  114. Self {
  115. dataset_path,
  116. nickname,
  117. workspace,
  118. workspaces,
  119. notify_queue_sender,
  120. write,
  121. password,
  122. p2p,
  123. event_graph,
  124. rpc_connections: Mutex::new(HashSet::new()),
  125. dnet_sub,
  126. deg_sub,
  127. }
  128. }
  129. // RPCAPI:
  130. // Activate or deactivate dnet in the P2P stack.
  131. // By sending `true`, dnet will be activated, and by sending `false` dnet will
  132. // be deactivated. Returns `true` on success.
  133. //
  134. // --> {"jsonrpc": "2.0", "method": "dnet_switch", "params": [true], "id": 42}
  135. // <-- {"jsonrpc": "2.0", "result": true, "id": 42}
  136. async fn dnet_switch(&self, params: JsonValue) -> TaudResult<JsonValue> {
  137. let params = params.get::<Vec<JsonValue>>().unwrap();
  138. if params.len() != 1 || !params[0].is_bool() {
  139. return Err(TaudError::InvalidData("Invalid parameters".into()))
  140. }
  141. let switch = params[0].get::<bool>().unwrap();
  142. if *switch {
  143. self.p2p.dnet_enable().await;
  144. } else {
  145. self.p2p.dnet_disable().await;
  146. }
  147. Ok(JsonValue::Boolean(true))
  148. }
  149. // RPCAPI:
  150. // Initializes a subscription to p2p dnet events.
  151. // Once a subscription is established, `darkirc` will send JSON-RPC notifications of
  152. // new network events to the subscriber.
  153. //
  154. // --> {"jsonrpc": "2.0", "method": "dnet.subscribe_events", "params": [], "id": 1}
  155. // <-- {"jsonrpc": "2.0", "method": "dnet.subscribe_events", "params": [`event`]}
  156. pub async fn dnet_subscribe_events(&self, id: u16, params: JsonValue) -> JsonResult {
  157. let params = params.get::<Vec<JsonValue>>().unwrap();
  158. if !params.is_empty() {
  159. return JsonError::new(ErrorCode::InvalidParams, None, id).into()
  160. }
  161. self.dnet_sub.clone().into()
  162. }
  163. // RPCAPI:
  164. // Initializes a subscription to deg events.
  165. // Once a subscription is established, apps using eventgraph will send JSON-RPC notifications of
  166. // new eventgraph events to the subscriber.
  167. //
  168. // --> {"jsonrpc": "2.0", "method": "deg.subscribe_events", "params": [], "id": 1}
  169. // <-- {"jsonrpc": "2.0", "method": "deg.subscribe_events", "params": [`event`]}
  170. pub async fn deg_subscribe_events(&self, id: u16, params: JsonValue) -> JsonResult {
  171. let params = params.get::<Vec<JsonValue>>().unwrap();
  172. if !params.is_empty() {
  173. return JsonError::new(ErrorCode::InvalidParams, None, id).into()
  174. }
  175. self.deg_sub.clone().into()
  176. }
  177. // RPCAPI:
  178. // Activate or deactivate deg in the EVENTGRAPH.
  179. // By sending `true`, deg will be activated, and by sending `false` deg
  180. // will be deactivated. Returns `true` on success.
  181. //
  182. // --> {"jsonrpc": "2.0", "method": "deg.switch", "params": [true], "id": 42}
  183. // <-- {"jsonrpc": "2.0", "result": true, "id": 42}
  184. async fn deg_switch(&self, _id: u16, params: JsonValue) -> TaudResult<JsonValue> {
  185. let params = params.get::<Vec<JsonValue>>().unwrap();
  186. if params.len() != 1 || !params[0].is_bool() {
  187. return Err(TaudError::InvalidData("Invalid parameters".into()))
  188. }
  189. let switch = params[0].get::<bool>().unwrap();
  190. if *switch {
  191. self.event_graph.deg_enable().await;
  192. } else {
  193. self.event_graph.deg_disable().await;
  194. }
  195. Ok(JsonValue::Boolean(true))
  196. }
  197. // RPCAPI:
  198. // Get EVENTGRAPH info.
  199. //
  200. // --> {"jsonrpc": "2.0", "method": "deg.switch", "params": [true], "id": 42}
  201. // <-- {"jsonrpc": "2.0", "result": true, "id": 42}
  202. async fn eg_get_info(&self, id: u16, params: JsonValue) -> JsonResult {
  203. let params_ = params.get::<Vec<JsonValue>>().unwrap();
  204. if !params_.is_empty() {
  205. return JsonError::new(ErrorCode::InvalidParams, None, id).into()
  206. }
  207. self.event_graph.eventgraph_info(id, params).await
  208. }
  209. // RPCAPI:
  210. // Add new task and returns `true` upon success.
  211. // --> {"jsonrpc": "2.0", "method": "add",
  212. // "params":
  213. // [{
  214. // "title": "..",
  215. // "desc": "..",
  216. // assign: [..],
  217. // project: [..],
  218. // "due": ..,
  219. // "rank": ..
  220. // }],
  221. // "id": 1
  222. // }
  223. // <-- {"jsonrpc": "2.0", "result": true, "id": 1}
  224. async fn add(&self, params: JsonValue) -> TaudResult<JsonValue> {
  225. let params = params.get::<Vec<JsonValue>>().unwrap();
  226. debug!(target: "tau", "JsonRpc::add() params {:?}", params);
  227. if !params[0].is_object() {
  228. return Err(TaudError::InvalidData("Invalid parameters".to_string()))
  229. }
  230. let params = params[0].get::<HashMap<String, JsonValue>>().unwrap();
  231. if params.len() != 9 {
  232. return Err(TaudError::InvalidData("Invalid parameters".to_string()))
  233. }
  234. let due = match params["due"] {
  235. JsonValue::Null => None,
  236. JsonValue::Number(numba) => Some(Timestamp::from_u64(numba as u64)),
  237. _ => return Err(TaudError::InvalidData("Invalid parameter \"due\"".to_string())),
  238. };
  239. let rank = match params["rank"] {
  240. JsonValue::Null => None,
  241. JsonValue::Number(numba) => Some(numba as f32),
  242. _ => return Err(TaudError::InvalidData("Invalid parameter \"rank\"".to_string())),
  243. };
  244. let tags = {
  245. let mut tags = vec![];
  246. for val in params["tags"].get::<Vec<JsonValue>>().unwrap().iter() {
  247. if let Some(tag) = val.get::<String>() {
  248. tags.push(tag.clone());
  249. } else {
  250. return Err(TaudError::InvalidData("Invalid parameter \"tags\"".to_string()))
  251. }
  252. }
  253. tags
  254. };
  255. let assigns = {
  256. let mut assigns = vec![];
  257. for val in params["assign"].get::<Vec<JsonValue>>().unwrap().iter() {
  258. if let Some(assign) = val.get::<String>() {
  259. assigns.push(assign.clone());
  260. } else {
  261. return Err(TaudError::InvalidData("Invalid parameter \"assign\"".to_string()))
  262. }
  263. }
  264. assigns
  265. };
  266. let projects = {
  267. let mut projects = vec![];
  268. for val in params["project"].get::<Vec<JsonValue>>().unwrap().iter() {
  269. if let Some(project) = val.get::<String>() {
  270. projects.push(project.clone());
  271. } else {
  272. return Err(TaudError::InvalidData("Invalid parameter \"project\"".to_string()))
  273. }
  274. }
  275. projects
  276. };
  277. let created_at = match params["created_at"] {
  278. JsonValue::Number(numba) => Some(numba as u64),
  279. _ => return Err(TaudError::InvalidData("Invalid parameter \"created_at\"".to_string())),
  280. };
  281. if !check_write_access(self.write.clone(), self.password.clone())? {
  282. return Ok(JsonValue::Boolean(false))
  283. }
  284. let mut new_task: TaskInfo = TaskInfo::new(
  285. self.workspace.lock().await.clone(),
  286. params["title"].get::<String>().unwrap(),
  287. params["desc"].get::<String>().unwrap(),
  288. &self.nickname,
  289. due,
  290. rank,
  291. Timestamp::from_u64(created_at.unwrap()),
  292. )?;
  293. new_task.set_project(&projects);
  294. new_task.set_assign(&assigns);
  295. new_task.set_tags(&tags);
  296. self.notify_queue_sender.send(new_task.clone()).await.map_err(Error::from)?;
  297. Ok(JsonValue::Boolean(true))
  298. }
  299. // RPCAPI:
  300. // List tasks
  301. // --> {"jsonrpc": "2.0", "method": "get_ids", "params": [], "id": 1}
  302. // <-- {"jsonrpc": "2.0", "result": [task_id, ...], "id": 1}
  303. async fn get_ref_ids(&self, params: JsonValue) -> TaudResult<JsonValue> {
  304. let params = params.get::<Vec<JsonValue>>().unwrap();
  305. debug!(target: "tau", "JsonRpc::get_ids() params {:?}", params);
  306. let ws = self.workspace.lock().await.clone();
  307. let tasks = MonthTasks::load_current_tasks(&self.dataset_path, ws, false)?;
  308. let task_ref_ids: Vec<JsonValue> =
  309. tasks.iter().map(|task| JsonValue::String(task.get_ref_id())).collect();
  310. Ok(JsonValue::Array(task_ref_ids))
  311. }
  312. // RPCAPI:
  313. // List tasks
  314. // --> {"jsonrpc": "2.0", "method": "get_ids", "params": [], "id": 1}
  315. // <-- {"jsonrpc": "2.0", "result": [task_id, ...], "id": 1}
  316. async fn get_archive_ref_ids(&self, params: JsonValue) -> TaudResult<JsonValue> {
  317. let params = params.get::<Vec<JsonValue>>().unwrap();
  318. debug!(target: "tau", "JsonRpc::get_archive_ref_ids() params {:?}", params);
  319. let month = match params[0].get::<String>() {
  320. Some(u64_str) => match u64_str.parse::<u64>() {
  321. Ok(v) => Some(Timestamp::from_u64(v)),
  322. //Err(e) => return Err(TaudError::InvalidData(e.to_string())),
  323. Err(_) => None,
  324. },
  325. None => None,
  326. };
  327. let ws = self.workspace.lock().await.clone();
  328. let tasks = MonthTasks::load_stop_tasks(&self.dataset_path, ws, month.as_ref())?;
  329. let task_ref_ids: Vec<JsonValue> =
  330. tasks.iter().map(|task| JsonValue::String(task.get_ref_id())).collect();
  331. Ok(JsonValue::Array(task_ref_ids))
  332. }
  333. // RPCAPI:
  334. // Modify task and returns `true` upon success.
  335. // --> {"jsonrpc": "2.0", "method": "modify", "params": [task_id, {"title": "new title"} ], "id": 1}
  336. // <-- {"jsonrpc": "2.0", "result": true, "id": 1}
  337. async fn modify(&self, params: JsonValue) -> TaudResult<JsonValue> {
  338. let params = params.get::<Vec<JsonValue>>().unwrap();
  339. debug!(target: "tau", "JsonRpc::modify() params {:?}", params);
  340. if params.len() != 2 || !params[0].is_string() || !params[1].is_object() {
  341. return Err(TaudError::InvalidData("len of params should be 2".into()))
  342. }
  343. if !check_write_access(self.write.clone(), self.password.clone())? {
  344. return Ok(JsonValue::Boolean(false))
  345. }
  346. let ws = self.workspace.lock().await.clone();
  347. let task = self.check_params_for_modify(
  348. params[0].get::<String>().unwrap(),
  349. params[1].get::<HashMap<String, JsonValue>>().unwrap(),
  350. ws,
  351. )?;
  352. self.notify_queue_sender.send(task).await.map_err(Error::from)?;
  353. Ok(JsonValue::Boolean(true))
  354. }
  355. // RPCAPI:
  356. // Set state for a task and returns `true` upon success.
  357. // --> {"jsonrpc": "2.0", "method": "set_state", "params": [task_id, state], "id": 1}
  358. // <-- {"jsonrpc": "2.0", "result": true, "id": 1}
  359. async fn set_state(&self, params: JsonValue) -> TaudResult<JsonValue> {
  360. // Allowed states for a task
  361. let states = ["stop", "start", "open", "pause"];
  362. let params = params.get::<Vec<JsonValue>>().unwrap();
  363. debug!(target: "tau", "JsonRpc::set_state() params {:?}", params);
  364. if params.len() != 2 || !params[0].is_string() || !params[1].is_string() {
  365. return Err(TaudError::InvalidData("len of params should be 2".into()))
  366. }
  367. if !check_write_access(self.write.clone(), self.password.clone())? {
  368. return Ok(JsonValue::Boolean(false))
  369. }
  370. let state = params[1].get::<String>().unwrap();
  371. let ws = self.workspace.lock().await.clone();
  372. let mut task: TaskInfo =
  373. self.load_task_by_ref_id(params[0].get::<String>().unwrap(), ws)?;
  374. if states.contains(&state.as_str()) {
  375. task.set_state(state);
  376. set_event(&mut task, "state", &self.nickname, state);
  377. }
  378. self.notify_queue_sender.send(task).await.map_err(Error::from)?;
  379. Ok(JsonValue::Boolean(true))
  380. }
  381. // RPCAPI:
  382. // Set comment for a task and returns `true` upon success.
  383. // --> {"jsonrpc": "2.0", "method": "set_comment", "params": [task_id, comment_content], "id": 1}
  384. // <-- {"jsonrpc": "2.0", "result": true, "id": 1}
  385. async fn set_comment(&self, params: JsonValue) -> TaudResult<JsonValue> {
  386. let params = params.get::<Vec<JsonValue>>().unwrap();
  387. debug!(target: "tau", "JsonRpc::set_comment() params {:?}", params);
  388. if params.len() != 2 || !params[0].is_string() || !params[1].is_string() {
  389. return Err(TaudError::InvalidData("len of params should be 2".into()))
  390. }
  391. if !check_write_access(self.write.clone(), self.password.clone())? {
  392. return Ok(JsonValue::Boolean(false))
  393. }
  394. let ref_id = params[0].get::<String>().unwrap();
  395. let comment_content = params[1].get::<String>().unwrap();
  396. let ws = self.workspace.lock().await.clone();
  397. let mut task: TaskInfo = self.load_task_by_ref_id(ref_id, ws)?;
  398. task.set_comment(Comment::new(comment_content, &self.nickname));
  399. set_event(&mut task, "comment", &self.nickname, comment_content);
  400. self.notify_queue_sender.send(task).await.map_err(Error::from)?;
  401. Ok(JsonValue::Boolean(true))
  402. }
  403. // RPCAPI:
  404. // Get a task by id.
  405. // --> {"jsonrpc": "2.0", "method": "get_task_by_id", "params": [task_id], "id": 1}
  406. // <-- {"jsonrpc": "2.0", "result": "task", "id": 1}
  407. async fn get_task_by_ref_id(&self, params: JsonValue) -> TaudResult<JsonValue> {
  408. let params = params.get::<Vec<JsonValue>>().unwrap();
  409. debug!(target: "tau", "JsonRpc::get_task_by_ref_id() params {:?}", params);
  410. if params.len() != 1 || !params[0].is_string() {
  411. return Err(TaudError::InvalidData("len of params should be 1".into()))
  412. }
  413. let ws = self.workspace.lock().await.clone();
  414. let task: TaskInfo = self.load_task_by_ref_id(params[0].get::<String>().unwrap(), ws)?;
  415. let task: JsonValue = (&task).into();
  416. Ok(task)
  417. }
  418. // RPCAPI:
  419. // Get all tasks.
  420. // --> {"jsonrpc": "2.0", "method": "fetch_deactive_tasks", "params": [task_id], "id": 1}
  421. // <-- {"jsonrpc": "2.0", "result": "task", "id": 1}
  422. async fn fetch_deactive_tasks(&self, params: JsonValue) -> TaudResult<JsonValue> {
  423. let params = params.get::<Vec<JsonValue>>().unwrap();
  424. debug!(target: "tau", "JsonRpc::fetch_deactive_tasks() params {:?}", params);
  425. if params.len() != 1 || !params[0].is_string() {
  426. return Err(TaudError::InvalidData("len of params should be 1".into()))
  427. }
  428. let month = match params[0].get::<String>() {
  429. Some(u64_str) => match u64_str.parse::<u64>() {
  430. Ok(v) => Some(Timestamp::from_u64(v)),
  431. //Err(e) => return Err(TaudError::InvalidData(e.to_string())),
  432. Err(_) => None,
  433. },
  434. None => None,
  435. };
  436. let ws = self.workspace.lock().await.clone();
  437. let tasks = MonthTasks::load_stop_tasks(&self.dataset_path, ws, month.as_ref())?;
  438. let tasks: Vec<JsonValue> = tasks.iter().map(|x| x.into()).collect();
  439. Ok(JsonValue::Array(tasks))
  440. }
  441. async fn fetch_archive_task(&self, params: JsonValue) -> TaudResult<JsonValue> {
  442. let params = params.get::<Vec<JsonValue>>().unwrap();
  443. debug!(target: "tau", "JsonRpc::fetch_archive_task() params {:?}", params);
  444. if params.len() != 2 || !params[0].is_string() || !params[1].is_string() {
  445. return Err(TaudError::InvalidData("len of params should be 2".into()))
  446. }
  447. let ref_id = params[0].get::<String>().unwrap();
  448. let month = match params[1].get::<String>() {
  449. Some(u64_str) => match u64_str.parse::<u64>() {
  450. Ok(v) => Some(Timestamp::from_u64(v)),
  451. //Err(e) => return Err(TaudError::InvalidData(e.to_string())),
  452. Err(_) => None,
  453. },
  454. None => None,
  455. };
  456. let ws = self.workspace.lock().await.clone();
  457. let mut tasks = MonthTasks::load_stop_tasks(&self.dataset_path, ws, month.as_ref())?;
  458. tasks.retain(|x| x.ref_id == *ref_id);
  459. if tasks.len() != 1 {
  460. return Err(TaudError::InvalidData("Must return a single value".into()))
  461. }
  462. let task: JsonValue = (&tasks[0]).into();
  463. Ok(task)
  464. }
  465. // RPCAPI:
  466. // Switch tasks workspace.
  467. // --> {"jsonrpc": "2.0", "method": "switch_ws", "params": [workspace], "id": 1}
  468. // <-- {"jsonrpc": "2.0", "result": "true", "id": 1}
  469. async fn switch_ws(&self, params: JsonValue) -> TaudResult<JsonValue> {
  470. let params = params.get::<Vec<JsonValue>>().unwrap();
  471. debug!(target: "tau", "JsonRpc::switch_ws() params {:?}", params);
  472. if params.len() != 1 {
  473. return Err(TaudError::InvalidData("len of params should be 1".into()))
  474. }
  475. if !params[0].is_string() {
  476. return Err(TaudError::InvalidData("Invalid workspace".into()))
  477. }
  478. let ws = params[0].get::<String>().unwrap();
  479. let mut s = self.workspace.lock().await;
  480. if self.workspaces.contains_key(ws) {
  481. *s = ws.to_string()
  482. } else {
  483. warn!("Workspace \"{}\" is not configured", ws);
  484. return Ok(JsonValue::Boolean(false))
  485. }
  486. Ok(JsonValue::Boolean(true))
  487. }
  488. // RPCAPI:
  489. // Get workspace.
  490. // --> {"jsonrpc": "2.0", "method": "get_ws", "params": [], "id": 1}
  491. // <-- {"jsonrpc": "2.0", "result": "workspace", "id": 1}
  492. async fn get_ws(&self, params: JsonValue) -> TaudResult<JsonValue> {
  493. let params = params.get::<Vec<JsonValue>>().unwrap();
  494. debug!(target: "tau", "JsonRpc::get_ws() params {:?}", params);
  495. let ws = self.workspace.lock().await.clone();
  496. Ok(JsonValue::String(ws))
  497. }
  498. // RPCAPI:
  499. // Export tasks.
  500. // --> {"jsonrpc": "2.0", "method": "export_to", "params": [path], "id": 1}
  501. // <-- {"jsonrpc": "2.0", "result": "true", "id": 1}
  502. async fn export_to(&self, params: JsonValue) -> TaudResult<JsonValue> {
  503. let params = params.get::<Vec<JsonValue>>().unwrap();
  504. debug!(target: "tau", "JsonRpc::export_to() params {:?}", params);
  505. if params.len() != 1 {
  506. return Err(TaudError::InvalidData("len of params should be 1".into()))
  507. }
  508. if !params[0].is_string() {
  509. return Err(TaudError::InvalidData("Invalid path".into()))
  510. }
  511. // mkdir datastore_path if not exists
  512. let path = params[0].get::<String>().unwrap();
  513. let path = expand_path(path)?.join("exported_tasks");
  514. create_dir_all(path.join("month")).map_err(Error::from)?;
  515. create_dir_all(path.join("task")).map_err(Error::from)?;
  516. let ws = self.workspace.lock().await.clone();
  517. let tasks = MonthTasks::load_current_tasks(&self.dataset_path, ws, true)?;
  518. for task in tasks {
  519. task.save(&path)?;
  520. }
  521. Ok(JsonValue::Boolean(true))
  522. }
  523. // RPCAPI:
  524. // Import tasks.
  525. // --> {"jsonrpc": "2.0", "method": "import_from", "params": [path], "id": 1}
  526. // <-- {"jsonrpc": "2.0", "result": "true", "id": 1}
  527. async fn import_from(&self, params: JsonValue) -> TaudResult<JsonValue> {
  528. let params = params.get::<Vec<JsonValue>>().unwrap();
  529. debug!(target: "tau", "JsonRpc::import_from() params {:?}", params);
  530. if params.len() != 1 {
  531. return Err(TaudError::InvalidData("len of params should be 1".into()))
  532. }
  533. if !params[0].is_string() {
  534. return Err(TaudError::InvalidData("Invalid path".into()))
  535. }
  536. if !check_write_access(self.write.clone(), self.password.clone())? {
  537. return Ok(JsonValue::Boolean(false))
  538. }
  539. let path = params[0].get::<String>().unwrap();
  540. let path = expand_path(path)?.join("exported_tasks");
  541. let ws = self.workspace.lock().await.clone();
  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. }