app_controller.rs 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290
  1. use crate::{
  2. entities::{
  3. app::{App, CreateAppParams, *},
  4. trash::TrashType,
  5. },
  6. errors::*,
  7. module::{WorkspaceDatabase, WorkspaceUser},
  8. notify::*,
  9. services::{helper::spawn, server::Server, TrashCan, TrashEvent},
  10. sql_tables::app::{AppTable, AppTableChangeset, AppTableSql},
  11. };
  12. use flowy_database::SqliteConnection;
  13. use futures::{FutureExt, StreamExt};
  14. use std::{collections::HashSet, sync::Arc};
  15. pub(crate) struct AppController {
  16. user: Arc<dyn WorkspaceUser>,
  17. database: Arc<dyn WorkspaceDatabase>,
  18. trash_can: Arc<TrashCan>,
  19. server: Server,
  20. }
  21. impl AppController {
  22. pub(crate) fn new(
  23. user: Arc<dyn WorkspaceUser>,
  24. database: Arc<dyn WorkspaceDatabase>,
  25. trash_can: Arc<TrashCan>,
  26. server: Server,
  27. ) -> Self {
  28. Self {
  29. user,
  30. database,
  31. trash_can,
  32. server,
  33. }
  34. }
  35. pub fn init(&self) -> Result<(), WorkspaceError> {
  36. self.listen_trash_can_event();
  37. Ok(())
  38. }
  39. #[tracing::instrument(level = "debug", skip(self), err)]
  40. pub(crate) async fn create_app(&self, params: CreateAppParams) -> Result<App, WorkspaceError> {
  41. let app = self.create_app_on_server(params).await?;
  42. let conn = &*self.database.db_connection()?;
  43. conn.immediate_transaction::<_, WorkspaceError, _>(|| {
  44. let _ = self.save_app(app.clone(), &*conn)?;
  45. let _ = notify_apps_changed(&app.workspace_id, self.trash_can.clone(), conn)?;
  46. Ok(())
  47. })?;
  48. Ok(app)
  49. }
  50. pub(crate) fn save_app(&self, app: App, conn: &SqliteConnection) -> Result<(), WorkspaceError> {
  51. let app_table = AppTable::new(app.clone());
  52. let _ = AppTableSql::create_app(app_table, &*conn)?;
  53. Ok(())
  54. }
  55. pub(crate) async fn read_app(&self, params: AppIdentifier) -> Result<App, WorkspaceError> {
  56. let conn = self.database.db_connection()?;
  57. let app_table = AppTableSql::read_app(&params.app_id, &*conn)?;
  58. let trash_ids = self.trash_can.trash_ids(&conn)?;
  59. if trash_ids.contains(&app_table.id) {
  60. return Err(WorkspaceError::record_not_found());
  61. }
  62. let _ = self.read_app_on_server(params)?;
  63. Ok(app_table.into())
  64. }
  65. pub(crate) async fn update_app(&self, params: UpdateAppParams) -> Result<(), WorkspaceError> {
  66. let changeset = AppTableChangeset::new(params.clone());
  67. let app_id = changeset.id.clone();
  68. let conn = &*self.database.db_connection()?;
  69. conn.immediate_transaction::<_, WorkspaceError, _>(|| {
  70. let _ = AppTableSql::update_app(changeset, conn)?;
  71. let app: App = AppTableSql::read_app(&app_id, conn)?.into();
  72. send_dart_notification(&app_id, WorkspaceNotification::AppUpdated)
  73. .payload(app)
  74. .send();
  75. Ok(())
  76. })?;
  77. let _ = self.update_app_on_server(params)?;
  78. Ok(())
  79. }
  80. pub(crate) fn read_app_tables(&self, ids: Vec<String>) -> Result<Vec<AppTable>, WorkspaceError> {
  81. let conn = &*self.database.db_connection()?;
  82. let mut app_tables = vec![];
  83. conn.immediate_transaction::<_, WorkspaceError, _>(|| {
  84. for app_id in ids {
  85. app_tables.push(AppTableSql::read_app(&app_id, conn)?);
  86. }
  87. Ok(())
  88. })?;
  89. Ok(app_tables)
  90. }
  91. }
  92. impl AppController {
  93. #[tracing::instrument(level = "debug", skip(self), err)]
  94. async fn create_app_on_server(&self, params: CreateAppParams) -> Result<App, WorkspaceError> {
  95. let token = self.user.token()?;
  96. let app = self.server.create_app(&token, params).await?;
  97. Ok(app)
  98. }
  99. #[tracing::instrument(level = "debug", skip(self), err)]
  100. fn update_app_on_server(&self, params: UpdateAppParams) -> Result<(), WorkspaceError> {
  101. let token = self.user.token()?;
  102. let server = self.server.clone();
  103. spawn(async move {
  104. match server.update_app(&token, params).await {
  105. Ok(_) => {},
  106. Err(e) => {
  107. // TODO: retry?
  108. log::error!("Update app failed: {:?}", e);
  109. },
  110. }
  111. });
  112. Ok(())
  113. }
  114. #[tracing::instrument(level = "debug", skip(self), err)]
  115. fn read_app_on_server(&self, params: AppIdentifier) -> Result<(), WorkspaceError> {
  116. let token = self.user.token()?;
  117. let server = self.server.clone();
  118. let pool = self.database.db_pool()?;
  119. spawn(async move {
  120. // Opti: retry?
  121. match server.read_app(&token, params).await {
  122. Ok(Some(app)) => match pool.get() {
  123. Ok(conn) => {
  124. let app_table = AppTable::new(app.clone());
  125. let result = AppTableSql::create_app(app_table, &*conn);
  126. match result {
  127. Ok(_) => {
  128. send_dart_notification(&app.id, WorkspaceNotification::AppUpdated)
  129. .payload(app)
  130. .send();
  131. },
  132. Err(e) => log::error!("Save app failed: {:?}", e),
  133. }
  134. },
  135. Err(e) => log::error!("Require db connection failed: {:?}", e),
  136. },
  137. Ok(None) => {},
  138. Err(e) => log::error!("Read app failed: {:?}", e),
  139. }
  140. });
  141. Ok(())
  142. }
  143. fn listen_trash_can_event(&self) {
  144. let mut rx = self.trash_can.subscribe();
  145. let database = self.database.clone();
  146. let trash_can = self.trash_can.clone();
  147. let _ = tokio::spawn(async move {
  148. loop {
  149. let mut stream = Box::pin(rx.recv().into_stream().filter_map(|result| async move {
  150. match result {
  151. Ok(event) => event.select(TrashType::App),
  152. Err(_e) => None,
  153. }
  154. }));
  155. match stream.next().await {
  156. Some(event) => handle_trash_event(database.clone(), trash_can.clone(), event).await,
  157. None => {},
  158. }
  159. }
  160. });
  161. }
  162. }
  163. #[tracing::instrument(level = "debug", skip(database, trash_can))]
  164. async fn handle_trash_event(database: Arc<dyn WorkspaceDatabase>, trash_can: Arc<TrashCan>, event: TrashEvent) {
  165. let db_result = database.db_connection();
  166. match event {
  167. TrashEvent::NewTrash(identifiers, ret) | TrashEvent::Putback(identifiers, ret) => {
  168. let result = || {
  169. let conn = &*db_result?;
  170. let _ = conn.immediate_transaction::<_, WorkspaceError, _>(|| {
  171. for identifier in identifiers.items {
  172. let app_table = AppTableSql::read_app(&identifier.id, conn)?;
  173. let _ = notify_apps_changed(&app_table.workspace_id, trash_can.clone(), conn)?;
  174. }
  175. Ok(())
  176. })?;
  177. Ok::<(), WorkspaceError>(())
  178. };
  179. let _ = ret.send(result()).await;
  180. },
  181. TrashEvent::Delete(identifiers, ret) => {
  182. let result = || {
  183. let conn = &*db_result?;
  184. let _ = conn.immediate_transaction::<_, WorkspaceError, _>(|| {
  185. let mut notify_ids = HashSet::new();
  186. for identifier in identifiers.items {
  187. let app_table = AppTableSql::read_app(&identifier.id, conn)?;
  188. let _ = AppTableSql::delete_app(&identifier.id, conn)?;
  189. notify_ids.insert(app_table.workspace_id);
  190. }
  191. for notify_id in notify_ids {
  192. let _ = notify_apps_changed(&notify_id, trash_can.clone(), conn)?;
  193. }
  194. Ok(())
  195. })?;
  196. Ok::<(), WorkspaceError>(())
  197. };
  198. let _ = ret.send(result()).await;
  199. },
  200. }
  201. }
  202. #[tracing::instrument(skip(workspace_id, trash_can, conn), err)]
  203. fn notify_apps_changed(workspace_id: &str, trash_can: Arc<TrashCan>, conn: &SqliteConnection) -> WorkspaceResult<()> {
  204. let repeated_app = read_local_workspace_apps(workspace_id, trash_can, conn)?;
  205. send_dart_notification(workspace_id, WorkspaceNotification::WorkspaceAppsChanged)
  206. .payload(repeated_app)
  207. .send();
  208. Ok(())
  209. }
  210. pub fn read_local_workspace_apps(
  211. workspace_id: &str,
  212. trash_can: Arc<TrashCan>,
  213. conn: &SqliteConnection,
  214. ) -> Result<RepeatedApp, WorkspaceError> {
  215. let mut app_tables = AppTableSql::read_workspace_apps(workspace_id, false, conn)?;
  216. let trash_ids = trash_can.trash_ids(conn)?;
  217. app_tables.retain(|app_table| !trash_ids.contains(&app_table.id));
  218. let apps = app_tables.into_iter().map(|table| table.into()).collect::<Vec<App>>();
  219. Ok(RepeatedApp { items: apps })
  220. }
  221. // #[tracing::instrument(level = "debug", skip(self), err)]
  222. // pub(crate) async fn delete_app(&self, app_id: &str) -> Result<(),
  223. // WorkspaceError> { let conn = &*self.database.db_connection()?;
  224. // conn.immediate_transaction::<_, WorkspaceError, _>(|| {
  225. // let app = AppTableSql::delete_app(app_id, &*conn)?;
  226. // let apps = self.read_local_apps(&app.workspace_id, &*conn)?;
  227. // send_dart_notification(&app.workspace_id,
  228. // WorkspaceNotification::WorkspaceDeleteApp) .payload(apps)
  229. // .send();
  230. // Ok(())
  231. // })?;
  232. //
  233. // let _ = self.delete_app_on_server(app_id);
  234. // Ok(())
  235. // }
  236. //
  237. // #[tracing::instrument(level = "debug", skip(self), err)]
  238. // fn delete_app_on_server(&self, app_id: &str) -> Result<(), WorkspaceError> {
  239. // let token = self.user.token()?;
  240. // let server = self.server.clone();
  241. // let params = DeleteAppParams {
  242. // app_id: app_id.to_string(),
  243. // };
  244. // spawn(async move {
  245. // match server.delete_app(&token, params).await {
  246. // Ok(_) => {},
  247. // Err(e) => {
  248. // // TODO: retry?
  249. // log::error!("Delete app failed: {:?}", e);
  250. // },
  251. // }
  252. // });
  253. // // let action = RetryAction::new(self.server.clone(), self.user.clone(),
  254. // move // |token, server| { let params = params.clone();
  255. // // async move {
  256. // // match server.delete_app(&token, params).await {
  257. // // Ok(_) => {},
  258. // // Err(e) => log::error!("Delete app failed: {:?}", e),
  259. // // }
  260. // // Ok::<(), WorkspaceError>(())
  261. // // }
  262. // // });
  263. // //
  264. // // spawn_retry(500, 3, action);
  265. // Ok(())
  266. // }