controller.rs 19 KB

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