controller.rs 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470
  1. use crate::manager::{ViewDataProcessor, ViewDataProcessorMap};
  2. use crate::{
  3. dart_notification::{send_dart_notification, FolderNotification},
  4. entities::{
  5. trash::{RepeatedTrashId, TrashType},
  6. view::{CreateViewParams, RepeatedView, UpdateViewParams, View, ViewId},
  7. },
  8. errors::{FlowyError, FlowyResult},
  9. event_map::{FolderCouldServiceV1, WorkspaceUser},
  10. services::{
  11. persistence::{FolderPersistence, FolderPersistenceTransaction, ViewChangeset},
  12. TrashController, TrashEvent,
  13. },
  14. };
  15. use bytes::Bytes;
  16. use flowy_database::kv::KV;
  17. use flowy_folder_data_model::entities::view::{gen_view_id, ViewDataType};
  18. use flowy_sync::entities::text_block_info::TextBlockId;
  19. use futures::{FutureExt, StreamExt};
  20. use std::{collections::HashSet, sync::Arc};
  21. const LATEST_VIEW_ID: &str = "latest_view_id";
  22. pub(crate) struct ViewController {
  23. user: Arc<dyn WorkspaceUser>,
  24. cloud_service: Arc<dyn FolderCouldServiceV1>,
  25. persistence: Arc<FolderPersistence>,
  26. trash_controller: Arc<TrashController>,
  27. data_processors: ViewDataProcessorMap,
  28. }
  29. impl ViewController {
  30. pub(crate) fn new(
  31. user: Arc<dyn WorkspaceUser>,
  32. persistence: Arc<FolderPersistence>,
  33. cloud_service: Arc<dyn FolderCouldServiceV1>,
  34. trash_controller: Arc<TrashController>,
  35. data_processors: ViewDataProcessorMap,
  36. ) -> Self {
  37. Self {
  38. user,
  39. cloud_service,
  40. persistence,
  41. trash_controller,
  42. data_processors,
  43. }
  44. }
  45. pub(crate) fn initialize(&self) -> Result<(), FlowyError> {
  46. self.listen_trash_can_event();
  47. Ok(())
  48. }
  49. #[tracing::instrument(level = "trace", skip(self, params), fields(name = %params.name), err)]
  50. pub(crate) async fn create_view_from_params(&self, mut params: CreateViewParams) -> Result<View, FlowyError> {
  51. let processor = self.get_data_processor(&params.data_type)?;
  52. let user_id = self.user.user_id()?;
  53. if params.data.is_empty() {
  54. let view_data = processor.create_default_view(&user_id, &params.view_id).await?;
  55. params.data = view_data.to_vec();
  56. } else {
  57. let delta_data = processor
  58. .process_view_delta_data(&user_id, &params.view_id, params.data.clone())
  59. .await?;
  60. let _ = self
  61. .create_view(&params.view_id, params.data_type.clone(), delta_data)
  62. .await?;
  63. };
  64. let view = self.create_view_on_server(params).await?;
  65. let _ = self.create_view_on_local(view.clone()).await?;
  66. Ok(view)
  67. }
  68. #[tracing::instrument(level = "debug", skip(self, view_id, delta_data), err)]
  69. pub(crate) async fn create_view(
  70. &self,
  71. view_id: &str,
  72. data_type: ViewDataType,
  73. delta_data: Bytes,
  74. ) -> Result<(), FlowyError> {
  75. if delta_data.is_empty() {
  76. return Err(FlowyError::internal().context("The content of the view should not be empty"));
  77. }
  78. let user_id = self.user.user_id()?;
  79. let processor = self.get_data_processor(&data_type)?;
  80. let _ = processor.create_container(&user_id, view_id, delta_data).await?;
  81. Ok(())
  82. }
  83. pub(crate) async fn create_view_on_local(&self, view: View) -> Result<(), FlowyError> {
  84. let trash_controller = self.trash_controller.clone();
  85. self.persistence
  86. .begin_transaction(|transaction| {
  87. let belong_to_id = view.belong_to_id.clone();
  88. let _ = transaction.create_view(view)?;
  89. let _ = notify_views_changed(&belong_to_id, trash_controller, &transaction)?;
  90. Ok(())
  91. })
  92. .await
  93. }
  94. #[tracing::instrument(level = "debug", skip(self, view_id), fields(view_id = %view_id.value), err)]
  95. pub(crate) async fn read_view(&self, view_id: ViewId) -> Result<View, FlowyError> {
  96. let view = self
  97. .persistence
  98. .begin_transaction(|transaction| {
  99. let view = transaction.read_view(&view_id.value)?;
  100. let trash_ids = self.trash_controller.read_trash_ids(&transaction)?;
  101. if trash_ids.contains(&view.id) {
  102. return Err(FlowyError::record_not_found());
  103. }
  104. Ok(view)
  105. })
  106. .await?;
  107. let _ = self.read_view_on_server(view_id);
  108. Ok(view)
  109. }
  110. pub(crate) async fn read_local_views(&self, ids: Vec<String>) -> Result<Vec<View>, FlowyError> {
  111. self.persistence
  112. .begin_transaction(|transaction| {
  113. let mut views = vec![];
  114. for view_id in ids {
  115. views.push(transaction.read_view(&view_id)?);
  116. }
  117. Ok(views)
  118. })
  119. .await
  120. }
  121. #[tracing::instrument(level = "trace", skip(self), err)]
  122. pub(crate) fn set_latest_view(&self, view_id: &str) -> Result<(), FlowyError> {
  123. KV::set_str(LATEST_VIEW_ID, view_id.to_owned());
  124. Ok(())
  125. }
  126. #[tracing::instrument(level = "debug", skip(self), err)]
  127. pub(crate) async fn close_view(&self, view_id: &str) -> Result<(), FlowyError> {
  128. let processor = self.get_data_processor_from_view_id(view_id).await?;
  129. let _ = processor.close_container(view_id).await?;
  130. Ok(())
  131. }
  132. #[tracing::instrument(level = "debug", skip(self,params), fields(doc_id = %params.value), err)]
  133. pub(crate) async fn delete_view(&self, params: TextBlockId) -> Result<(), FlowyError> {
  134. if let Some(view_id) = KV::get_str(LATEST_VIEW_ID) {
  135. if view_id == params.value {
  136. let _ = KV::remove(LATEST_VIEW_ID);
  137. }
  138. }
  139. let processor = self.get_data_processor_from_view_id(&params.value).await?;
  140. let _ = processor.delete_container(&params.value).await?;
  141. Ok(())
  142. }
  143. #[tracing::instrument(level = "debug", skip(self), err)]
  144. pub(crate) async fn move_view(&self, view_id: &str, from: usize, to: usize) -> Result<(), FlowyError> {
  145. let _ = self
  146. .persistence
  147. .begin_transaction(|transaction| {
  148. let _ = transaction.move_view(view_id, from, to)?;
  149. let view = transaction.read_view(view_id)?;
  150. let _ = notify_views_changed(&view.belong_to_id, self.trash_controller.clone(), &transaction)?;
  151. Ok(())
  152. })
  153. .await?;
  154. Ok(())
  155. }
  156. #[tracing::instrument(level = "debug", skip(self), err)]
  157. pub(crate) async fn duplicate_view(&self, view_id: &str) -> Result<(), FlowyError> {
  158. let view = self
  159. .persistence
  160. .begin_transaction(|transaction| transaction.read_view(view_id))
  161. .await?;
  162. let processor = self.get_data_processor(&view.data_type)?;
  163. let delta_bytes = processor.view_delta_data(view_id).await?;
  164. let duplicate_params = CreateViewParams {
  165. belong_to_id: view.belong_to_id.clone(),
  166. name: format!("{} (copy)", &view.name),
  167. desc: view.desc,
  168. thumbnail: view.thumbnail,
  169. data_type: view.data_type,
  170. data: delta_bytes.to_vec(),
  171. view_id: gen_view_id(),
  172. plugin_type: view.plugin_type,
  173. };
  174. let _ = self.create_view_from_params(duplicate_params).await?;
  175. Ok(())
  176. }
  177. // belong_to_id will be the app_id or view_id.
  178. #[tracing::instrument(level = "trace", skip(self), err)]
  179. pub(crate) async fn read_views_belong_to(&self, belong_to_id: &str) -> Result<RepeatedView, FlowyError> {
  180. self.persistence
  181. .begin_transaction(|transaction| {
  182. read_belonging_views_on_local(belong_to_id, self.trash_controller.clone(), &transaction)
  183. })
  184. .await
  185. }
  186. #[tracing::instrument(level = "debug", skip(self, params), err)]
  187. pub(crate) async fn update_view(&self, params: UpdateViewParams) -> Result<View, FlowyError> {
  188. let changeset = ViewChangeset::new(params.clone());
  189. let view_id = changeset.id.clone();
  190. let view = self
  191. .persistence
  192. .begin_transaction(|transaction| {
  193. let _ = transaction.update_view(changeset)?;
  194. let view = transaction.read_view(&view_id)?;
  195. send_dart_notification(&view_id, FolderNotification::ViewUpdated)
  196. .payload(view.clone())
  197. .send();
  198. let _ = notify_views_changed(&view.belong_to_id, self.trash_controller.clone(), &transaction)?;
  199. Ok(view)
  200. })
  201. .await?;
  202. let _ = self.update_view_on_server(params);
  203. Ok(view)
  204. }
  205. pub(crate) async fn latest_visit_view(&self) -> FlowyResult<Option<View>> {
  206. match KV::get_str(LATEST_VIEW_ID) {
  207. None => Ok(None),
  208. Some(view_id) => {
  209. let view = self
  210. .persistence
  211. .begin_transaction(|transaction| transaction.read_view(&view_id))
  212. .await?;
  213. Ok(Some(view))
  214. }
  215. }
  216. }
  217. }
  218. impl ViewController {
  219. #[tracing::instrument(level = "debug", skip(self, params), err)]
  220. async fn create_view_on_server(&self, params: CreateViewParams) -> Result<View, FlowyError> {
  221. let token = self.user.token()?;
  222. let view = self.cloud_service.create_view(&token, params).await?;
  223. Ok(view)
  224. }
  225. #[tracing::instrument(level = "debug", skip(self), err)]
  226. fn update_view_on_server(&self, params: UpdateViewParams) -> Result<(), FlowyError> {
  227. let token = self.user.token()?;
  228. let server = self.cloud_service.clone();
  229. tokio::spawn(async move {
  230. match server.update_view(&token, params).await {
  231. Ok(_) => {}
  232. Err(e) => {
  233. // TODO: retry?
  234. log::error!("Update view failed: {:?}", e);
  235. }
  236. }
  237. });
  238. Ok(())
  239. }
  240. #[tracing::instrument(level = "debug", skip(self), err)]
  241. fn read_view_on_server(&self, params: ViewId) -> Result<(), FlowyError> {
  242. let token = self.user.token()?;
  243. let server = self.cloud_service.clone();
  244. let persistence = self.persistence.clone();
  245. // TODO: Retry with RetryAction?
  246. tokio::spawn(async move {
  247. match server.read_view(&token, params).await {
  248. Ok(Some(view)) => {
  249. match persistence
  250. .begin_transaction(|transaction| transaction.create_view(view.clone()))
  251. .await
  252. {
  253. Ok(_) => {
  254. send_dart_notification(&view.id, FolderNotification::ViewUpdated)
  255. .payload(view.clone())
  256. .send();
  257. }
  258. Err(e) => log::error!("Save view failed: {:?}", e),
  259. }
  260. }
  261. Ok(None) => {}
  262. Err(e) => log::error!("Read view failed: {:?}", e),
  263. }
  264. });
  265. Ok(())
  266. }
  267. fn listen_trash_can_event(&self) {
  268. let mut rx = self.trash_controller.subscribe();
  269. let persistence = self.persistence.clone();
  270. let data_processors = self.data_processors.clone();
  271. let trash_controller = self.trash_controller.clone();
  272. let _ = tokio::spawn(async move {
  273. loop {
  274. let mut stream = Box::pin(rx.recv().into_stream().filter_map(|result| async move {
  275. match result {
  276. Ok(event) => event.select(TrashType::TrashView),
  277. Err(_e) => None,
  278. }
  279. }));
  280. if let Some(event) = stream.next().await {
  281. handle_trash_event(
  282. persistence.clone(),
  283. data_processors.clone(),
  284. trash_controller.clone(),
  285. event,
  286. )
  287. .await
  288. }
  289. }
  290. });
  291. }
  292. async fn get_data_processor_from_view_id(
  293. &self,
  294. view_id: &str,
  295. ) -> FlowyResult<Arc<dyn ViewDataProcessor + Send + Sync>> {
  296. let view = self
  297. .persistence
  298. .begin_transaction(|transaction| transaction.read_view(view_id))
  299. .await?;
  300. self.get_data_processor(&view.data_type)
  301. }
  302. #[inline]
  303. fn get_data_processor(&self, data_type: &ViewDataType) -> FlowyResult<Arc<dyn ViewDataProcessor + Send + Sync>> {
  304. match self.data_processors.get(data_type) {
  305. None => Err(FlowyError::internal().context(format!(
  306. "Get data processor failed. Unknown view data type: {:?}",
  307. data_type
  308. ))),
  309. Some(processor) => Ok(processor.clone()),
  310. }
  311. }
  312. }
  313. #[tracing::instrument(level = "trace", skip(persistence, data_processors, trash_can))]
  314. async fn handle_trash_event(
  315. persistence: Arc<FolderPersistence>,
  316. data_processors: ViewDataProcessorMap,
  317. trash_can: Arc<TrashController>,
  318. event: TrashEvent,
  319. ) {
  320. match event {
  321. TrashEvent::NewTrash(identifiers, ret) => {
  322. let result = persistence
  323. .begin_transaction(|transaction| {
  324. let views = read_local_views_with_transaction(identifiers, &transaction)?;
  325. for view in views {
  326. let _ = notify_views_changed(&view.belong_to_id, trash_can.clone(), &transaction)?;
  327. notify_dart(view, FolderNotification::ViewDeleted);
  328. }
  329. Ok(())
  330. })
  331. .await;
  332. let _ = ret.send(result).await;
  333. }
  334. TrashEvent::Putback(identifiers, ret) => {
  335. let result = persistence
  336. .begin_transaction(|transaction| {
  337. let views = read_local_views_with_transaction(identifiers, &transaction)?;
  338. for view in views {
  339. let _ = notify_views_changed(&view.belong_to_id, trash_can.clone(), &transaction)?;
  340. notify_dart(view, FolderNotification::ViewRestored);
  341. }
  342. Ok(())
  343. })
  344. .await;
  345. let _ = ret.send(result).await;
  346. }
  347. TrashEvent::Delete(identifiers, ret) => {
  348. let result = || async {
  349. let views = persistence
  350. .begin_transaction(|transaction| {
  351. let mut notify_ids = HashSet::new();
  352. let mut views = vec![];
  353. for identifier in identifiers.items {
  354. let view = transaction.read_view(&identifier.id)?;
  355. let _ = transaction.delete_view(&view.id)?;
  356. notify_ids.insert(view.belong_to_id.clone());
  357. views.push(view);
  358. }
  359. for notify_id in notify_ids {
  360. let _ = notify_views_changed(&notify_id, trash_can.clone(), &transaction)?;
  361. }
  362. Ok(views)
  363. })
  364. .await?;
  365. for view in views {
  366. match get_data_processor(data_processors.clone(), &view.data_type) {
  367. Ok(processor) => {
  368. let _ = processor.close_container(&view.id).await?;
  369. }
  370. Err(e) => {
  371. tracing::error!("{}", e)
  372. }
  373. }
  374. }
  375. Ok(())
  376. };
  377. let _ = ret.send(result().await).await;
  378. }
  379. }
  380. }
  381. fn get_data_processor(
  382. data_processors: ViewDataProcessorMap,
  383. data_type: &ViewDataType,
  384. ) -> FlowyResult<Arc<dyn ViewDataProcessor + Send + Sync>> {
  385. match data_processors.get(data_type) {
  386. None => Err(FlowyError::internal().context(format!(
  387. "Get data processor failed. Unknown view data type: {:?}",
  388. data_type
  389. ))),
  390. Some(processor) => Ok(processor.clone()),
  391. }
  392. }
  393. fn read_local_views_with_transaction<'a>(
  394. identifiers: RepeatedTrashId,
  395. transaction: &'a (dyn FolderPersistenceTransaction + 'a),
  396. ) -> Result<Vec<View>, FlowyError> {
  397. let mut views = vec![];
  398. for identifier in identifiers.items {
  399. let view = transaction.read_view(&identifier.id)?;
  400. views.push(view);
  401. }
  402. Ok(views)
  403. }
  404. fn notify_dart(view: View, notification: FolderNotification) {
  405. send_dart_notification(&view.id, notification).payload(view).send();
  406. }
  407. #[tracing::instrument(
  408. level = "debug",
  409. skip(belong_to_id, trash_controller, transaction),
  410. fields(view_count),
  411. err
  412. )]
  413. fn notify_views_changed<'a>(
  414. belong_to_id: &str,
  415. trash_controller: Arc<TrashController>,
  416. transaction: &'a (dyn FolderPersistenceTransaction + 'a),
  417. ) -> FlowyResult<()> {
  418. let repeated_view = read_belonging_views_on_local(belong_to_id, trash_controller.clone(), transaction)?;
  419. tracing::Span::current().record("view_count", &format!("{}", repeated_view.len()).as_str());
  420. send_dart_notification(belong_to_id, FolderNotification::AppViewsChanged)
  421. .payload(repeated_view)
  422. .send();
  423. Ok(())
  424. }
  425. fn read_belonging_views_on_local<'a>(
  426. belong_to_id: &str,
  427. trash_controller: Arc<TrashController>,
  428. transaction: &'a (dyn FolderPersistenceTransaction + 'a),
  429. ) -> FlowyResult<RepeatedView> {
  430. let mut views = transaction.read_views(belong_to_id)?;
  431. let trash_ids = trash_controller.read_trash_ids(transaction)?;
  432. views.retain(|view_table| !trash_ids.contains(&view_table.id));
  433. Ok(RepeatedView { items: views })
  434. }