controller.rs 20 KB

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