ws_actor.rs 4.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156
  1. use crate::{
  2. context::FlowyPersistence,
  3. services::web_socket::{entities::Socket, WSClientData, WSUser, WebSocketMessage},
  4. util::serde_ext::{md5, parse_from_bytes},
  5. };
  6. use actix_rt::task::spawn_blocking;
  7. use async_stream::stream;
  8. use backend_service::errors::{internal_error, Result, ServerError};
  9. use flowy_collaboration::{
  10. protobuf::{
  11. ClientRevisionWSData as ClientRevisionWSDataPB,
  12. ClientRevisionWSDataType as ClientRevisionWSDataTypePB,
  13. Revision as RevisionPB,
  14. },
  15. server_document::ServerDocumentManager,
  16. synchronizer::{RevisionSyncResponse, RevisionUser},
  17. };
  18. use futures::stream::StreamExt;
  19. use std::sync::Arc;
  20. use tokio::sync::{mpsc, oneshot};
  21. pub enum DocumentWSActorMessage {
  22. ClientData {
  23. client_data: WSClientData,
  24. persistence: Arc<FlowyPersistence>,
  25. ret: oneshot::Sender<Result<()>>,
  26. },
  27. }
  28. pub struct DocumentWebSocketActor {
  29. receiver: Option<mpsc::Receiver<DocumentWSActorMessage>>,
  30. doc_manager: Arc<ServerDocumentManager>,
  31. }
  32. impl DocumentWebSocketActor {
  33. pub fn new(receiver: mpsc::Receiver<DocumentWSActorMessage>, manager: Arc<ServerDocumentManager>) -> Self {
  34. Self {
  35. receiver: Some(receiver),
  36. doc_manager: manager,
  37. }
  38. }
  39. pub async fn run(mut self) {
  40. let mut receiver = self
  41. .receiver
  42. .take()
  43. .expect("DocumentWebSocketActor's receiver should only take one time");
  44. let stream = stream! {
  45. loop {
  46. match receiver.recv().await {
  47. Some(msg) => yield msg,
  48. None => break,
  49. }
  50. }
  51. };
  52. stream.for_each(|msg| self.handle_message(msg)).await;
  53. }
  54. async fn handle_message(&self, msg: DocumentWSActorMessage) {
  55. match msg {
  56. DocumentWSActorMessage::ClientData {
  57. client_data,
  58. persistence: _,
  59. ret,
  60. } => {
  61. let _ = ret.send(self.handle_document_data(client_data).await);
  62. },
  63. }
  64. }
  65. async fn handle_document_data(&self, client_data: WSClientData) -> Result<()> {
  66. let WSClientData { user, socket, data } = client_data;
  67. let document_client_data = spawn_blocking(move || parse_from_bytes::<ClientRevisionWSDataPB>(&data))
  68. .await
  69. .map_err(internal_error)??;
  70. tracing::debug!(
  71. "[DocumentWebSocketActor]: receive: {}:{}, {:?}",
  72. document_client_data.object_id,
  73. document_client_data.data_id,
  74. document_client_data.ty
  75. );
  76. let user = Arc::new(DocumentRevisionUser { user, socket });
  77. match &document_client_data.ty {
  78. ClientRevisionWSDataTypePB::ClientPushRev => {
  79. let _ = self
  80. .doc_manager
  81. .handle_client_revisions(user, document_client_data)
  82. .await
  83. .map_err(internal_error)?;
  84. },
  85. ClientRevisionWSDataTypePB::ClientPing => {
  86. let _ = self
  87. .doc_manager
  88. .handle_client_ping(user, document_client_data)
  89. .await
  90. .map_err(internal_error)?;
  91. },
  92. }
  93. Ok(())
  94. }
  95. }
  96. #[allow(dead_code)]
  97. fn verify_md5(revision: &RevisionPB) -> Result<()> {
  98. if md5(&revision.delta_data) != revision.md5 {
  99. return Err(ServerError::internal().context("RevisionPB md5 not match"));
  100. }
  101. Ok(())
  102. }
  103. #[derive(Clone)]
  104. pub struct DocumentRevisionUser {
  105. pub user: Arc<WSUser>,
  106. pub(crate) socket: Socket,
  107. }
  108. impl std::fmt::Debug for DocumentRevisionUser {
  109. fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
  110. f.debug_struct("DocumentRevisionUser")
  111. .field("user", &self.user)
  112. .field("socket", &self.socket)
  113. .finish()
  114. }
  115. }
  116. impl RevisionUser for DocumentRevisionUser {
  117. fn user_id(&self) -> String { self.user.id().to_string() }
  118. fn receive(&self, resp: RevisionSyncResponse) {
  119. let result = match resp {
  120. RevisionSyncResponse::Pull(data) => {
  121. let msg: WebSocketMessage = data.into();
  122. self.socket.try_send(msg).map_err(internal_error)
  123. },
  124. RevisionSyncResponse::Push(data) => {
  125. let msg: WebSocketMessage = data.into();
  126. self.socket.try_send(msg).map_err(internal_error)
  127. },
  128. RevisionSyncResponse::Ack(data) => {
  129. let msg: WebSocketMessage = data.into();
  130. self.socket.try_send(msg).map_err(internal_error)
  131. },
  132. };
  133. match result {
  134. Ok(_) => {},
  135. Err(e) => log::error!("[DocumentRevisionUser]: {}", e),
  136. }
  137. }
  138. }