controller.rs 17 KB

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