edit_doc.rs 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329
  1. use crate::{
  2. entities::{
  3. doc::{DocDelta, RevId, RevType, Revision, RevisionRange},
  4. ws::{WsDataType, WsDocumentData},
  5. },
  6. errors::{internal_error, DocError, DocResult},
  7. module::DocumentUser,
  8. services::{
  9. doc::{
  10. edit::{
  11. doc_actor::DocumentActor,
  12. message::{DocumentMsg, TransformDeltas},
  13. model::OpenDocAction,
  14. },
  15. revision::{RevisionManager, RevisionServer},
  16. UndoResult,
  17. },
  18. ws::{DocumentWebSocket, WsDocumentHandler},
  19. },
  20. };
  21. use bytes::Bytes;
  22. use flowy_database::ConnectionPool;
  23. use flowy_infra::retry::{ExponentialBackoff, Retry};
  24. use flowy_ot::core::{Attribute, Delta, Interval};
  25. use flowy_ws::WsState;
  26. use std::{convert::TryFrom, sync::Arc};
  27. use tokio::sync::{mpsc, mpsc::UnboundedSender, oneshot};
  28. pub type DocId = String;
  29. pub struct ClientEditDoc {
  30. pub doc_id: DocId,
  31. rev_manager: Arc<RevisionManager>,
  32. document: UnboundedSender<DocumentMsg>,
  33. ws: Arc<dyn DocumentWebSocket>,
  34. user: Arc<dyn DocumentUser>,
  35. }
  36. impl ClientEditDoc {
  37. pub(crate) async fn new(
  38. doc_id: &str,
  39. pool: Arc<ConnectionPool>,
  40. ws: Arc<dyn DocumentWebSocket>,
  41. server: Arc<dyn RevisionServer>,
  42. user: Arc<dyn DocumentUser>,
  43. ) -> DocResult<Self> {
  44. let (sender, receiver) = mpsc::unbounded_channel();
  45. let mut rev_manager = RevisionManager::new(doc_id, pool.clone(), server.clone(), sender);
  46. spawn_rev_receiver(receiver, ws.clone());
  47. let delta = rev_manager.load_document().await?;
  48. let document = spawn_doc_edit_actor(doc_id, delta, pool.clone());
  49. let doc_id = doc_id.to_string();
  50. let rev_manager = Arc::new(rev_manager);
  51. let edit_doc = Self {
  52. doc_id,
  53. rev_manager,
  54. document,
  55. ws,
  56. user,
  57. };
  58. edit_doc.notify_open_doc();
  59. Ok(edit_doc)
  60. }
  61. pub async fn insert<T: ToString>(&self, index: usize, data: T) -> Result<(), DocError> {
  62. let (ret, rx) = oneshot::channel::<DocResult<Delta>>();
  63. let msg = DocumentMsg::Insert {
  64. index,
  65. data: data.to_string(),
  66. ret,
  67. };
  68. let _ = self.document.send(msg);
  69. let delta = rx.await.map_err(internal_error)??;
  70. let rev_id = self.save_revision(delta).await?;
  71. save_document(self.document.clone(), rev_id.into()).await
  72. }
  73. pub async fn delete(&self, interval: Interval) -> Result<(), DocError> {
  74. let (ret, rx) = oneshot::channel::<DocResult<Delta>>();
  75. let msg = DocumentMsg::Delete { interval, ret };
  76. let _ = self.document.send(msg);
  77. let delta = rx.await.map_err(internal_error)??;
  78. let _ = self.save_revision(delta).await?;
  79. Ok(())
  80. }
  81. pub async fn format(&self, interval: Interval, attribute: Attribute) -> Result<(), DocError> {
  82. let (ret, rx) = oneshot::channel::<DocResult<Delta>>();
  83. let msg = DocumentMsg::Format {
  84. interval,
  85. attribute,
  86. ret,
  87. };
  88. let _ = self.document.send(msg);
  89. let delta = rx.await.map_err(internal_error)??;
  90. let _ = self.save_revision(delta).await?;
  91. Ok(())
  92. }
  93. pub async fn replace<T: ToString>(&mut self, interval: Interval, data: T) -> Result<(), DocError> {
  94. let (ret, rx) = oneshot::channel::<DocResult<Delta>>();
  95. let msg = DocumentMsg::Replace {
  96. interval,
  97. data: data.to_string(),
  98. ret,
  99. };
  100. let _ = self.document.send(msg);
  101. let delta = rx.await.map_err(internal_error)??;
  102. let _ = self.save_revision(delta).await?;
  103. Ok(())
  104. }
  105. pub async fn can_undo(&self) -> bool {
  106. let (ret, rx) = oneshot::channel::<bool>();
  107. let msg = DocumentMsg::CanUndo { ret };
  108. let _ = self.document.send(msg);
  109. rx.await.unwrap_or(false)
  110. }
  111. pub async fn can_redo(&self) -> bool {
  112. let (ret, rx) = oneshot::channel::<bool>();
  113. let msg = DocumentMsg::CanRedo { ret };
  114. let _ = self.document.send(msg);
  115. rx.await.unwrap_or(false)
  116. }
  117. pub async fn undo(&self) -> Result<UndoResult, DocError> {
  118. let (ret, rx) = oneshot::channel::<DocResult<UndoResult>>();
  119. let msg = DocumentMsg::Undo { ret };
  120. let _ = self.document.send(msg);
  121. rx.await.map_err(internal_error)?
  122. }
  123. pub async fn redo(&self) -> Result<UndoResult, DocError> {
  124. let (ret, rx) = oneshot::channel::<DocResult<UndoResult>>();
  125. let msg = DocumentMsg::Redo { ret };
  126. let _ = self.document.send(msg);
  127. rx.await.map_err(internal_error)?
  128. }
  129. pub async fn delta(&self) -> DocResult<DocDelta> {
  130. let (ret, rx) = oneshot::channel::<DocResult<String>>();
  131. let msg = DocumentMsg::Doc { ret };
  132. let _ = self.document.send(msg);
  133. let data = rx.await.map_err(internal_error)??;
  134. Ok(DocDelta {
  135. doc_id: self.doc_id.clone(),
  136. data,
  137. })
  138. }
  139. #[tracing::instrument(level = "debug", skip(self, delta), fields(revision_delta = %delta.to_json(), send_state, base_rev_id, rev_id))]
  140. async fn save_revision(&self, delta: Delta) -> Result<RevId, DocError> {
  141. let delta_data = delta.to_bytes();
  142. let (base_rev_id, rev_id) = self.rev_manager.next_rev_id();
  143. tracing::Span::current().record("base_rev_id", &base_rev_id);
  144. tracing::Span::current().record("rev_id", &rev_id);
  145. let delta_data = delta_data.to_vec();
  146. let revision = Revision::new(base_rev_id, rev_id, delta_data, &self.doc_id, RevType::Local);
  147. let _ = self.rev_manager.add_revision(&revision).await?;
  148. Ok(rev_id.into())
  149. }
  150. #[tracing::instrument(level = "debug", skip(self, data), err)]
  151. pub(crate) async fn compose_local_delta(&self, data: Bytes) -> Result<(), DocError> {
  152. let delta = Delta::from_bytes(&data)?;
  153. let (ret, rx) = oneshot::channel::<DocResult<()>>();
  154. let msg = DocumentMsg::Delta {
  155. delta: delta.clone(),
  156. ret,
  157. };
  158. let _ = self.document.send(msg);
  159. let _ = rx.await.map_err(internal_error)??;
  160. let rev_id = self.save_revision(delta).await?;
  161. save_document(self.document.clone(), rev_id).await
  162. }
  163. #[cfg(feature = "flowy_test")]
  164. pub async fn doc_json(&self) -> DocResult<String> {
  165. let (ret, rx) = oneshot::channel::<DocResult<String>>();
  166. let msg = DocumentMsg::Doc { ret };
  167. let _ = self.document.send(msg);
  168. rx.await.map_err(internal_error)?
  169. }
  170. #[tracing::instrument(level = "debug", skip(self))]
  171. fn notify_open_doc(&self) {
  172. let rev_id: RevId = self.rev_manager.rev_id().into();
  173. if let Ok(user_id) = self.user.user_id() {
  174. let action = OpenDocAction::new(&user_id, &self.doc_id, &rev_id, &self.ws);
  175. let strategy = ExponentialBackoff::from_millis(50).take(3);
  176. let retry = Retry::spawn(strategy, action);
  177. tokio::spawn(async move {
  178. match retry.await {
  179. Ok(_) => {},
  180. Err(e) => log::error!("Notify open doc failed: {}", e),
  181. }
  182. });
  183. }
  184. }
  185. #[tracing::instrument(level = "debug", skip(self))]
  186. async fn handle_push_rev(&self, bytes: Bytes) -> DocResult<()> {
  187. // Transform the revision
  188. let (ret, rx) = oneshot::channel::<DocResult<TransformDeltas>>();
  189. let _ = self.document.send(DocumentMsg::RemoteRevision { bytes, ret });
  190. let TransformDeltas {
  191. client_prime,
  192. server_prime,
  193. server_rev_id,
  194. } = rx.await.map_err(internal_error)??;
  195. if self.rev_manager.rev_id() >= server_rev_id.value {
  196. // Ignore this push revision if local_rev_id >= server_rev_id
  197. return Ok(());
  198. }
  199. // compose delta
  200. let (ret, rx) = oneshot::channel::<DocResult<()>>();
  201. let msg = DocumentMsg::Delta {
  202. delta: client_prime.clone(),
  203. ret,
  204. };
  205. let _ = self.document.send(msg);
  206. let _ = rx.await.map_err(internal_error)??;
  207. // update rev id
  208. self.rev_manager.set_rev_id(server_rev_id.clone().into());
  209. let (local_base_rev_id, local_rev_id) = self.rev_manager.next_rev_id();
  210. // save the revision
  211. let revision = Revision::new(
  212. local_base_rev_id,
  213. local_rev_id,
  214. client_prime.to_bytes().to_vec(),
  215. &self.doc_id,
  216. RevType::Remote,
  217. );
  218. let _ = self.rev_manager.add_revision(&revision).await?;
  219. // send the server_prime delta
  220. let revision = Revision::new(
  221. local_base_rev_id,
  222. local_rev_id,
  223. server_prime.to_bytes().to_vec(),
  224. &self.doc_id,
  225. RevType::Remote,
  226. );
  227. let _ = self.ws.send(revision.into());
  228. let _ = save_document(self.document.clone(), local_rev_id.into()).await?;
  229. Ok(())
  230. }
  231. async fn handle_ws_message(&self, doc_data: WsDocumentData) -> DocResult<()> {
  232. let bytes = Bytes::from(doc_data.data);
  233. match doc_data.ty {
  234. WsDataType::PushRev => {
  235. let _ = self.handle_push_rev(bytes).await?;
  236. },
  237. WsDataType::PullRev => {
  238. let range = RevisionRange::try_from(bytes)?;
  239. let revision = self.rev_manager.construct_revisions(range).await?;
  240. let _ = self.ws.send(revision.into());
  241. },
  242. WsDataType::NewDocUser => {},
  243. WsDataType::Acked => {
  244. let rev_id = RevId::try_from(bytes)?;
  245. let _ = self.rev_manager.ack_rev(rev_id).await?;
  246. },
  247. WsDataType::Conflict => {},
  248. }
  249. Ok(())
  250. }
  251. }
  252. pub struct EditDocWsHandler(pub Arc<ClientEditDoc>);
  253. impl WsDocumentHandler for EditDocWsHandler {
  254. fn receive(&self, doc_data: WsDocumentData) {
  255. let edit_doc = self.0.clone();
  256. tokio::spawn(async move {
  257. if let Err(e) = edit_doc.handle_ws_message(doc_data).await {
  258. log::error!("{:?}", e);
  259. }
  260. });
  261. }
  262. fn state_changed(&self, state: &WsState) {
  263. match state {
  264. WsState::Init => {},
  265. WsState::Connected(_) => self.0.notify_open_doc(),
  266. WsState::Disconnected(_e) => {},
  267. }
  268. }
  269. }
  270. fn spawn_rev_receiver(mut receiver: mpsc::UnboundedReceiver<Revision>, ws: Arc<dyn DocumentWebSocket>) {
  271. tokio::spawn(async move {
  272. loop {
  273. while let Some(revision) = receiver.recv().await {
  274. log::debug!("Send revision:{} to server", revision.rev_id);
  275. match ws.send(revision.into()) {
  276. Ok(_) => {},
  277. Err(e) => log::error!("Send revision failed: {:?}", e),
  278. };
  279. }
  280. }
  281. });
  282. }
  283. async fn save_document(document: UnboundedSender<DocumentMsg>, rev_id: RevId) -> DocResult<()> {
  284. let (ret, rx) = oneshot::channel::<DocResult<()>>();
  285. let _ = document.send(DocumentMsg::SaveDocument { rev_id, ret });
  286. let result = rx.await.map_err(internal_error)?;
  287. result
  288. }
  289. fn spawn_doc_edit_actor(_doc_id: &str, delta: Delta, _pool: Arc<ConnectionPool>) -> UnboundedSender<DocumentMsg> {
  290. let (sender, receiver) = mpsc::unbounded_channel::<DocumentMsg>();
  291. let actor = DocumentActor::new(delta, receiver);
  292. tokio::spawn(actor.run());
  293. sender
  294. }