controller.rs 17 KB

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