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 log::{debug, info, warn};
  26. use smol::lock::{Mutex, MutexGuard};
  27. use tinyjson::JsonValue;
  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. // TODO: make this optional
  84. "p2p.get_info" => return self.p2p_get_info(req.id, req.params).await,
  85. _ => return JsonError::new(ErrorCode::MethodNotFound, None, req.id).into(),
  86. };
  87. to_json_result(rep, req.id)
  88. }
  89. async fn connections_mut(&self) -> MutexGuard<'life0, HashSet<StoppableTaskPtr>> {
  90. self.rpc_connections.lock().await
  91. }
  92. }
  93. impl HandlerP2p for JsonRpcInterface {
  94. fn p2p(&self) -> net::P2pPtr {
  95. self.p2p.clone()
  96. }
  97. }
  98. impl JsonRpcInterface {
  99. #[allow(clippy::too_many_arguments)]
  100. pub fn new(
  101. dataset_path: PathBuf,
  102. notify_queue_sender: smol::channel::Sender<TaskInfo>,
  103. nickname: String,
  104. workspaces: Arc<HashMap<String, Workspace>>,
  105. p2p: net::P2pPtr,
  106. event_graph: EventGraphPtr,
  107. dnet_sub: JsonSubscriber,
  108. deg_sub: JsonSubscriber,
  109. ) -> Self {
  110. let workspace = Mutex::new(DEFAULT_WORKSPACE.to_string());
  111. Self {
  112. dataset_path,
  113. nickname,
  114. workspace,
  115. workspaces,
  116. notify_queue_sender,
  117. p2p,
  118. event_graph,
  119. rpc_connections: Mutex::new(HashSet::new()),
  120. dnet_sub,
  121. deg_sub,
  122. }
  123. }
  124. // RPCAPI:
  125. // Activate or deactivate dnet in the P2P stack.
  126. // By sending `true`, dnet will be activated, and by sending `false` dnet will
  127. // be deactivated. Returns `true` on success.
  128. //
  129. // --> {"jsonrpc": "2.0", "method": "dnet_switch", "params": [true], "id": 42}
  130. // <-- {"jsonrpc": "2.0", "result": true, "id": 42}
  131. async fn dnet_switch(&self, params: JsonValue) -> TaudResult<JsonValue> {
  132. let params = params.get::<Vec<JsonValue>>().unwrap();
  133. if params.len() != 1 || !params[0].is_bool() {
  134. return Err(TaudError::InvalidData("Invalid parameters".into()))
  135. }
  136. let switch = params[0].get::<bool>().unwrap();
  137. if *switch {
  138. self.p2p.dnet_enable();
  139. } else {
  140. self.p2p.dnet_disable();
  141. }
  142. Ok(JsonValue::Boolean(true))
  143. }
  144. // RPCAPI:
  145. // Initializes a subscription to p2p dnet events.
  146. // Once a subscription is established, `darkirc` will send JSON-RPC notifications of
  147. // new network events to the subscriber.
  148. //
  149. // --> {"jsonrpc": "2.0", "method": "dnet.subscribe_events", "params": [], "id": 1}
  150. // <-- {"jsonrpc": "2.0", "method": "dnet.subscribe_events", "params": [`event`]}
  151. pub async fn dnet_subscribe_events(&self, id: u16, params: JsonValue) -> JsonResult {
  152. let params = params.get::<Vec<JsonValue>>().unwrap();
  153. if !params.is_empty() {
  154. return JsonError::new(ErrorCode::InvalidParams, None, id).into()
  155. }
  156. self.dnet_sub.clone().into()
  157. }
  158. // RPCAPI:
  159. // Initializes a subscription to deg events.
  160. // Once a subscription is established, apps using eventgraph will send JSON-RPC notifications of
  161. // new eventgraph events to the subscriber.
  162. //
  163. // --> {"jsonrpc": "2.0", "method": "deg.subscribe_events", "params": [], "id": 1}
  164. // <-- {"jsonrpc": "2.0", "method": "deg.subscribe_events", "params": [`event`]}
  165. pub async fn deg_subscribe_events(&self, id: u16, params: JsonValue) -> JsonResult {
  166. let params = params.get::<Vec<JsonValue>>().unwrap();
  167. if !params.is_empty() {
  168. return JsonError::new(ErrorCode::InvalidParams, None, id).into()
  169. }
  170. self.deg_sub.clone().into()
  171. }
  172. // RPCAPI:
  173. // Activate or deactivate deg in the EVENTGRAPH.
  174. // By sending `true`, deg will be activated, and by sending `false` deg
  175. // will be deactivated. Returns `true` on success.
  176. //
  177. // --> {"jsonrpc": "2.0", "method": "deg.switch", "params": [true], "id": 42}
  178. // <-- {"jsonrpc": "2.0", "result": true, "id": 42}
  179. async fn deg_switch(&self, _id: u16, params: JsonValue) -> TaudResult<JsonValue> {
  180. let params = params.get::<Vec<JsonValue>>().unwrap();
  181. if params.len() != 1 || !params[0].is_bool() {
  182. return Err(TaudError::InvalidData("Invalid parameters".into()))
  183. }
  184. let switch = params[0].get::<bool>().unwrap();
  185. if *switch {
  186. self.event_graph.deg_enable().await;
  187. } else {
  188. self.event_graph.deg_disable().await;
  189. }
  190. Ok(JsonValue::Boolean(true))
  191. }
  192. // RPCAPI:
  193. // Get EVENTGRAPH info.
  194. //
  195. // --> {"jsonrpc": "2.0", "method": "deg.switch", "params": [true], "id": 42}
  196. // <-- {"jsonrpc": "2.0", "result": true, "id": 42}
  197. async fn eg_get_info(&self, id: u16, params: JsonValue) -> JsonResult {
  198. let params_ = params.get::<Vec<JsonValue>>().unwrap();
  199. if !params_.is_empty() {
  200. return JsonError::new(ErrorCode::InvalidParams, None, id).into()
  201. }
  202. self.event_graph.eventgraph_info(id, params).await
  203. }
  204. // RPCAPI:
  205. // Add new task and returns `true` upon success.
  206. // --> {"jsonrpc": "2.0", "method": "add",
  207. // "params":
  208. // [{
  209. // "title": "..",
  210. // "desc": "..",
  211. // assign: [..],
  212. // project: [..],
  213. // "due": ..,
  214. // "rank": ..
  215. // }],
  216. // "id": 1
  217. // }
  218. // <-- {"jsonrpc": "2.0", "result": true, "id": 1}
  219. async fn add(&self, params: JsonValue) -> TaudResult<JsonValue> {
  220. let params = params.get::<Vec<JsonValue>>().unwrap();
  221. debug!(target: "tau", "JsonRpc::add() params {:?}", params);
  222. if !params[0].is_object() {
  223. return Err(TaudError::InvalidData("Invalid parameters".to_string()))
  224. }
  225. let params = params[0].get::<HashMap<String, JsonValue>>().unwrap();
  226. if params.len() != 9 {
  227. return Err(TaudError::InvalidData("Invalid parameters".to_string()))
  228. }
  229. let due = match params["due"] {
  230. JsonValue::Null => None,
  231. JsonValue::Number(numba) => Some(Timestamp::from_u64(numba as u64)),
  232. _ => return Err(TaudError::InvalidData("Invalid parameter \"due\"".to_string())),
  233. };
  234. let rank = match params["rank"] {
  235. JsonValue::Null => None,
  236. JsonValue::Number(numba) => Some(numba as f32),
  237. _ => return Err(TaudError::InvalidData("Invalid parameter \"rank\"".to_string())),
  238. };
  239. let tags = {
  240. let mut tags = vec![];
  241. for val in params["tags"].get::<Vec<JsonValue>>().unwrap().iter() {
  242. if let Some(tag) = val.get::<String>() {
  243. tags.push(tag.clone());
  244. } else {
  245. return Err(TaudError::InvalidData("Invalid parameter \"tags\"".to_string()))
  246. }
  247. }
  248. tags
  249. };
  250. let assigns = {
  251. let mut assigns = vec![];
  252. for val in params["assign"].get::<Vec<JsonValue>>().unwrap().iter() {
  253. if let Some(assign) = val.get::<String>() {
  254. assigns.push(assign.clone());
  255. } else {
  256. return Err(TaudError::InvalidData("Invalid parameter \"assign\"".to_string()))
  257. }
  258. }
  259. assigns
  260. };
  261. let projects = {
  262. let mut projects = vec![];
  263. for val in params["project"].get::<Vec<JsonValue>>().unwrap().iter() {
  264. if let Some(project) = val.get::<String>() {
  265. projects.push(project.clone());
  266. } else {
  267. return Err(TaudError::InvalidData("Invalid parameter \"project\"".to_string()))
  268. }
  269. }
  270. projects
  271. };
  272. let created_at = match params["created_at"] {
  273. JsonValue::Number(numba) => Some(numba as u64),
  274. _ => return Err(TaudError::InvalidData("Invalid parameter \"created_at\"".to_string())),
  275. };
  276. let ws = self.workspace.lock().await.clone();
  277. if self.workspaces.get(&ws).unwrap().write_key.is_none() {
  278. info!("You don't have write access!");
  279. return Ok(JsonValue::Boolean(false))
  280. }
  281. let mut new_task: TaskInfo = TaskInfo::new(
  282. ws,
  283. params["title"].get::<String>().unwrap(),
  284. params["desc"].get::<String>().unwrap(),
  285. &self.nickname,
  286. due,
  287. rank,
  288. Timestamp::from_u64(created_at.unwrap()),
  289. )?;
  290. new_task.set_project(&projects);
  291. new_task.set_assign(&assigns);
  292. new_task.set_tags(&tags);
  293. self.notify_queue_sender.send(new_task.clone()).await.map_err(Error::from)?;
  294. Ok(JsonValue::Boolean(true))
  295. }
  296. // RPCAPI:
  297. // List tasks
  298. // --> {"jsonrpc": "2.0", "method": "get_ids", "params": [], "id": 1}
  299. // <-- {"jsonrpc": "2.0", "result": [task_id, ...], "id": 1}
  300. async fn get_ref_ids(&self, params: JsonValue) -> TaudResult<JsonValue> {
  301. let params = params.get::<Vec<JsonValue>>().unwrap();
  302. debug!(target: "tau", "JsonRpc::get_ids() params {:?}", params);
  303. let ws = self.workspace.lock().await.clone();
  304. let tasks = MonthTasks::load_current_tasks(&self.dataset_path, ws, false)?;
  305. let task_ref_ids: Vec<JsonValue> =
  306. tasks.iter().map(|task| JsonValue::String(task.get_ref_id())).collect();
  307. Ok(JsonValue::Array(task_ref_ids))
  308. }
  309. // RPCAPI:
  310. // List tasks
  311. // --> {"jsonrpc": "2.0", "method": "get_ids", "params": [], "id": 1}
  312. // <-- {"jsonrpc": "2.0", "result": [task_id, ...], "id": 1}
  313. async fn get_archive_ref_ids(&self, params: JsonValue) -> TaudResult<JsonValue> {
  314. let params = params.get::<Vec<JsonValue>>().unwrap();
  315. debug!(target: "tau", "JsonRpc::get_archive_ref_ids() params {:?}", params);
  316. let month = match params[0].get::<String>() {
  317. Some(u64_str) => match u64_str.parse::<u64>() {
  318. Ok(v) => Some(Timestamp::from_u64(v)),
  319. //Err(e) => return Err(TaudError::InvalidData(e.to_string())),
  320. Err(_) => None,
  321. },
  322. None => None,
  323. };
  324. let ws = self.workspace.lock().await.clone();
  325. let tasks = MonthTasks::load_stop_tasks(&self.dataset_path, ws, month.as_ref())?;
  326. let task_ref_ids: Vec<JsonValue> =
  327. tasks.iter().map(|task| JsonValue::String(task.get_ref_id())).collect();
  328. Ok(JsonValue::Array(task_ref_ids))
  329. }
  330. // RPCAPI:
  331. // Modify task and returns `true` upon success.
  332. // --> {"jsonrpc": "2.0", "method": "modify", "params": [task_id, {"title": "new title"} ], "id": 1}
  333. // <-- {"jsonrpc": "2.0", "result": true, "id": 1}
  334. async fn modify(&self, params: JsonValue) -> TaudResult<JsonValue> {
  335. let params = params.get::<Vec<JsonValue>>().unwrap();
  336. debug!(target: "tau", "JsonRpc::modify() params {:?}", params);
  337. if params.len() != 2 || !params[0].is_string() || !params[1].is_object() {
  338. return Err(TaudError::InvalidData("len of params should be 2".into()))
  339. }
  340. let ws = self.workspace.lock().await.clone();
  341. if self.workspaces.get(&ws).unwrap().write_key.is_none() {
  342. info!("You don't have write access!");
  343. return Ok(JsonValue::Boolean(false))
  344. }
  345. let task = self.check_params_for_modify(
  346. params[0].get::<String>().unwrap(),
  347. params[1].get::<HashMap<String, JsonValue>>().unwrap(),
  348. ws,
  349. )?;
  350. self.notify_queue_sender.send(task).await.map_err(Error::from)?;
  351. Ok(JsonValue::Boolean(true))
  352. }
  353. // RPCAPI:
  354. // Set state for a task and returns `true` upon success.
  355. // --> {"jsonrpc": "2.0", "method": "set_state", "params": [task_id, state], "id": 1}
  356. // <-- {"jsonrpc": "2.0", "result": true, "id": 1}
  357. async fn set_state(&self, params: JsonValue) -> TaudResult<JsonValue> {
  358. // Allowed states for a task
  359. let states = ["stop", "start", "open", "pause"];
  360. let params = params.get::<Vec<JsonValue>>().unwrap();
  361. debug!(target: "tau", "JsonRpc::set_state() params {:?}", params);
  362. if params.len() != 2 || !params[0].is_string() || !params[1].is_string() {
  363. return Err(TaudError::InvalidData("len of params should be 2".into()))
  364. }
  365. let state = params[1].get::<String>().unwrap();
  366. let ws = self.workspace.lock().await.clone();
  367. if self.workspaces.get(&ws).unwrap().write_key.is_none() {
  368. info!("You don't have write access!");
  369. return Ok(JsonValue::Boolean(false))
  370. }
  371. let mut task: TaskInfo =
  372. self.load_task_by_ref_id(params[0].get::<String>().unwrap(), ws)?;
  373. if states.contains(&state.as_str()) {
  374. task.set_state(state);
  375. set_event(&mut task, "state", &self.nickname, state);
  376. }
  377. self.notify_queue_sender.send(task).await.map_err(Error::from)?;
  378. Ok(JsonValue::Boolean(true))
  379. }
  380. // RPCAPI:
  381. // Set comment for a task and returns `true` upon success.
  382. // --> {"jsonrpc": "2.0", "method": "set_comment", "params": [task_id, comment_content], "id": 1}
  383. // <-- {"jsonrpc": "2.0", "result": true, "id": 1}
  384. async fn set_comment(&self, params: JsonValue) -> TaudResult<JsonValue> {
  385. let params = params.get::<Vec<JsonValue>>().unwrap();
  386. debug!(target: "tau", "JsonRpc::set_comment() params {:?}", params);
  387. if params.len() != 2 || !params[0].is_string() || !params[1].is_string() {
  388. return Err(TaudError::InvalidData("len of params should be 2".into()))
  389. }
  390. let ref_id = params[0].get::<String>().unwrap();
  391. let comment_content = params[1].get::<String>().unwrap();
  392. let ws = self.workspace.lock().await.clone();
  393. if self.workspaces.get(&ws).unwrap().write_key.is_none() {
  394. info!("You don't have write access!");
  395. return Ok(JsonValue::Boolean(false))
  396. }
  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. let path = params[0].get::<String>().unwrap();
  537. let path = expand_path(path)?.join("exported_tasks");
  538. let ws = self.workspace.lock().await.clone();
  539. if self.workspaces.get(&ws).unwrap().write_key.is_none() {
  540. info!("You don't have write access!");
  541. return Ok(JsonValue::Boolean(false))
  542. }
  543. let imported_tasks = MonthTasks::load_current_tasks(&path, ws.clone(), true)?;
  544. for task in imported_tasks {
  545. if MonthTasks::load_current_tasks(&self.dataset_path, ws.clone(), false)?
  546. .into_iter()
  547. .map(|t| t.ref_id)
  548. .any(|x| x == task.ref_id)
  549. {
  550. continue
  551. }
  552. self.notify_queue_sender.send(task).await.map_err(Error::from)?;
  553. }
  554. Ok(JsonValue::Boolean(true))
  555. }
  556. fn load_task_by_ref_id(&self, task_ref_id: &str, ws: String) -> TaudResult<TaskInfo> {
  557. let tasks = MonthTasks::load_current_tasks(&self.dataset_path, ws, false)?;
  558. let task = tasks.into_iter().find(|t| (t.get_ref_id()) == task_ref_id);
  559. task.ok_or(TaudError::InvalidId)
  560. }
  561. fn check_params_for_modify(
  562. &self,
  563. task_ref_id: &str,
  564. fields: &HashMap<String, JsonValue>,
  565. ws: String,
  566. ) -> TaudResult<TaskInfo> {
  567. let mut task: TaskInfo = self.load_task_by_ref_id(task_ref_id, ws)?;
  568. if fields.contains_key("title") {
  569. let title = fields["title"].get::<String>().unwrap();
  570. if !title.is_empty() {
  571. task.set_title(title);
  572. set_event(&mut task, "title", &self.nickname, title);
  573. }
  574. }
  575. if fields.contains_key("desc") {
  576. let desc = fields["desc"].get::<String>().unwrap();
  577. if !desc.is_empty() {
  578. task.set_desc(desc);
  579. set_event(&mut task, "desc", &self.nickname, desc);
  580. }
  581. }
  582. if fields.contains_key("rank") {
  583. match fields["rank"] {
  584. JsonValue::Null => set_event(&mut task, "rank", &self.nickname, "None"),
  585. JsonValue::Number(rank) => {
  586. task.set_rank(Some(rank as f32));
  587. set_event(&mut task, "rank", &self.nickname, &rank.to_string())
  588. }
  589. _ => unreachable!(),
  590. }
  591. }
  592. if fields.contains_key("due") {
  593. match &fields["due"] {
  594. JsonValue::Null => set_event(&mut task, "due", &self.nickname, "None"),
  595. JsonValue::Number(ts_num) => {
  596. task.set_due(Some(Timestamp::from_u64(*ts_num as u64)));
  597. set_event(&mut task, "due", &self.nickname, &ts_num.to_string())
  598. }
  599. _ => unreachable!(),
  600. }
  601. }
  602. if fields.contains_key("assign") {
  603. let assign: Vec<String> = fields["assign"]
  604. .get::<Vec<JsonValue>>()
  605. .unwrap()
  606. .iter()
  607. .map(|x| x.get::<String>().unwrap().clone())
  608. .collect();
  609. if !assign.is_empty() {
  610. task.set_assign(&assign);
  611. set_event(&mut task, "assign", &self.nickname, &assign.join(", "));
  612. }
  613. }
  614. if fields.contains_key("project") {
  615. let project: Vec<String> = fields["project"]
  616. .get::<Vec<JsonValue>>()
  617. .unwrap()
  618. .iter()
  619. .map(|x| x.get::<String>().unwrap().clone())
  620. .collect();
  621. if !project.is_empty() {
  622. task.set_project(&project);
  623. set_event(&mut task, "project", &self.nickname, &project.join(", "));
  624. }
  625. }
  626. if fields.contains_key("tags") {
  627. let tags: Vec<String> = fields["tags"]
  628. .get::<Vec<JsonValue>>()
  629. .unwrap()
  630. .iter()
  631. .map(|x| x.get::<String>().unwrap().clone())
  632. .collect();
  633. if !tags.is_empty() {
  634. task.set_tags(&tags);
  635. set_event(&mut task, "tags", &self.nickname, &tags.join(", "));
  636. }
  637. }
  638. Ok(task)
  639. }
  640. }