user_session.rs 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387
  1. use crate::{
  2. entities::{SignInParams, SignUpParams, UpdateUserParams, UserProfile},
  3. errors::{ErrorCode, UserError},
  4. services::user::database::UserDB,
  5. sql_tables::{UserTable, UserTableChangeset},
  6. };
  7. use crate::{
  8. notify::*,
  9. services::{
  10. server::{construct_user_server, Server},
  11. user::notifier::UserNotifier,
  12. },
  13. };
  14. use backend_service::config::ServerConfig;
  15. use flowy_database::{
  16. query_dsl::*,
  17. schema::{user_table, user_table::dsl},
  18. DBConnection,
  19. ExpressionMethods,
  20. UserDatabaseConnection,
  21. };
  22. use lib_infra::kv::KV;
  23. use lib_sqlite::ConnectionPool;
  24. use lib_ws::{WsController, WsMessageHandler, WsState};
  25. use parking_lot::RwLock;
  26. use serde::{Deserialize, Serialize};
  27. use std::sync::Arc;
  28. use tokio::sync::mpsc;
  29. pub struct UserSessionConfig {
  30. root_dir: String,
  31. server_config: ServerConfig,
  32. session_cache_key: String,
  33. }
  34. impl UserSessionConfig {
  35. pub fn new(root_dir: &str, server_config: &ServerConfig, session_cache_key: &str) -> Self {
  36. Self {
  37. root_dir: root_dir.to_owned(),
  38. server_config: server_config.clone(),
  39. session_cache_key: session_cache_key.to_owned(),
  40. }
  41. }
  42. }
  43. pub struct UserSession {
  44. database: UserDB,
  45. config: UserSessionConfig,
  46. #[allow(dead_code)]
  47. server: Server,
  48. session: RwLock<Option<Session>>,
  49. pub ws_controller: Arc<WsController>,
  50. pub notifier: UserNotifier,
  51. }
  52. impl UserSession {
  53. pub fn new(config: UserSessionConfig) -> Self {
  54. let db = UserDB::new(&config.root_dir);
  55. let server = construct_user_server(&config.server_config);
  56. let ws_controller = Arc::new(WsController::new());
  57. let notifier = UserNotifier::new();
  58. Self {
  59. database: db,
  60. config,
  61. server,
  62. session: RwLock::new(None),
  63. ws_controller,
  64. notifier,
  65. }
  66. }
  67. pub fn init(&self) {
  68. if let Ok(session) = self.get_session() {
  69. self.notifier.notify_login(&session.token);
  70. }
  71. }
  72. pub fn db_connection(&self) -> Result<DBConnection, UserError> {
  73. let user_id = self.get_session()?.user_id;
  74. self.database.get_connection(&user_id)
  75. }
  76. // The caller will be not 'Sync' before of the return value,
  77. // PooledConnection<ConnectionManager> is not sync. You can use
  78. // db_connection_pool function to require the ConnectionPool that is 'Sync'.
  79. //
  80. // let pool = self.db_connection_pool()?;
  81. // let conn: PooledConnection<ConnectionManager> = pool.get()?;
  82. pub fn db_pool(&self) -> Result<Arc<ConnectionPool>, UserError> {
  83. let user_id = self.get_session()?.user_id;
  84. self.database.get_pool(&user_id)
  85. }
  86. #[tracing::instrument(level = "debug", skip(self))]
  87. pub async fn sign_in(&self, params: SignInParams) -> Result<UserProfile, UserError> {
  88. if self.is_login(&params.email) {
  89. self.user_profile().await
  90. } else {
  91. let resp = self.server.sign_in(params).await?;
  92. let session = Session::new(&resp.user_id, &resp.token, &resp.email);
  93. let _ = self.set_session(Some(session))?;
  94. let user_table = self.save_user(resp.into()).await?;
  95. let user_profile: UserProfile = user_table.into();
  96. self.notifier.notify_login(&user_profile.token);
  97. Ok(user_profile)
  98. }
  99. }
  100. #[tracing::instrument(level = "debug", skip(self))]
  101. pub async fn sign_up(&self, params: SignUpParams) -> Result<UserProfile, UserError> {
  102. if self.is_login(&params.email) {
  103. self.user_profile().await
  104. } else {
  105. let resp = self.server.sign_up(params).await?;
  106. let session = Session::new(&resp.user_id, &resp.token, &resp.email);
  107. let _ = self.set_session(Some(session))?;
  108. let user_table = self.save_user(resp.into()).await?;
  109. let user_profile: UserProfile = user_table.into();
  110. let (ret, mut tx) = mpsc::channel(1);
  111. self.notifier.notify_sign_up(ret, &user_profile);
  112. let _ = tx.recv().await;
  113. Ok(user_profile)
  114. }
  115. }
  116. #[tracing::instrument(level = "debug", skip(self))]
  117. pub async fn sign_out(&self) -> Result<(), UserError> {
  118. let session = self.get_session()?;
  119. let _ =
  120. diesel::delete(dsl::user_table.filter(dsl::id.eq(&session.user_id))).execute(&*(self.db_connection()?))?;
  121. let _ = self.database.close_user_db(&session.user_id)?;
  122. let _ = self.set_session(None)?;
  123. self.notifier.notify_logout(&session.token);
  124. let _ = self.sign_out_on_server(&session.token).await?;
  125. Ok(())
  126. }
  127. #[tracing::instrument(level = "debug", skip(self))]
  128. pub async fn update_user(&self, params: UpdateUserParams) -> Result<(), UserError> {
  129. let session = self.get_session()?;
  130. let changeset = UserTableChangeset::new(params.clone());
  131. diesel_update_table!(user_table, changeset, &*self.db_connection()?);
  132. let _ = self.update_user_on_server(&session.token, params).await?;
  133. Ok(())
  134. }
  135. pub async fn init_user(&self) -> Result<(), UserError> {
  136. let (_, token) = self.get_session()?.into_part();
  137. let _ = self.start_ws_connection(&token).await?;
  138. Ok(())
  139. }
  140. pub async fn check_user(&self) -> Result<UserProfile, UserError> {
  141. let (user_id, token) = self.get_session()?.into_part();
  142. let user = dsl::user_table
  143. .filter(user_table::id.eq(&user_id))
  144. .first::<UserTable>(&*(self.db_connection()?))?;
  145. let _ = self.read_user_profile_on_server(&token)?;
  146. Ok(user.into())
  147. }
  148. pub async fn user_profile(&self) -> Result<UserProfile, UserError> {
  149. let (user_id, token) = self.get_session()?.into_part();
  150. let user = dsl::user_table
  151. .filter(user_table::id.eq(&user_id))
  152. .first::<UserTable>(&*(self.db_connection()?))?;
  153. let _ = self.read_user_profile_on_server(&token)?;
  154. Ok(user.into())
  155. }
  156. pub fn user_dir(&self) -> Result<String, UserError> {
  157. let session = self.get_session()?;
  158. Ok(format!("{}/{}", self.config.root_dir, session.user_id))
  159. }
  160. pub fn user_id(&self) -> Result<String, UserError> { Ok(self.get_session()?.user_id) }
  161. pub fn token(&self) -> Result<String, UserError> { Ok(self.get_session()?.token) }
  162. pub fn add_ws_handler(&self, handler: Arc<dyn WsMessageHandler>) {
  163. let _ = self.ws_controller.add_handler(handler);
  164. }
  165. pub fn update_network_state(&self, state: NetworkState) {
  166. log::info!("{:?}", state);
  167. }
  168. }
  169. impl UserSession {
  170. fn read_user_profile_on_server(&self, token: &str) -> Result<(), UserError> {
  171. let server = self.server.clone();
  172. let token = token.to_owned();
  173. tokio::spawn(async move {
  174. match server.get_user(&token).await {
  175. Ok(profile) => {
  176. dart_notify(&token, UserNotification::UserProfileUpdated)
  177. .payload(profile)
  178. .send();
  179. },
  180. Err(e) => {
  181. dart_notify(&token, UserNotification::UserProfileUpdated)
  182. .error(e)
  183. .send();
  184. },
  185. }
  186. });
  187. Ok(())
  188. }
  189. async fn update_user_on_server(&self, token: &str, params: UpdateUserParams) -> Result<(), UserError> {
  190. let server = self.server.clone();
  191. let token = token.to_owned();
  192. let _ = tokio::spawn(async move {
  193. match server.update_user(&token, params).await {
  194. Ok(_) => {},
  195. Err(e) => {
  196. // TODO: retry?
  197. log::error!("update user profile failed: {:?}", e);
  198. },
  199. }
  200. })
  201. .await;
  202. Ok(())
  203. }
  204. async fn sign_out_on_server(&self, token: &str) -> Result<(), UserError> {
  205. let server = self.server.clone();
  206. let token = token.to_owned();
  207. let _ = tokio::spawn(async move {
  208. match server.sign_out(&token).await {
  209. Ok(_) => {},
  210. Err(e) => log::error!("Sign out failed: {:?}", e),
  211. }
  212. })
  213. .await;
  214. Ok(())
  215. }
  216. async fn save_user(&self, user: UserTable) -> Result<UserTable, UserError> {
  217. let conn = self.db_connection()?;
  218. let _ = diesel::insert_into(user_table::table)
  219. .values(user.clone())
  220. .execute(&*conn)?;
  221. Ok(user)
  222. }
  223. fn set_session(&self, session: Option<Session>) -> Result<(), UserError> {
  224. tracing::debug!("Set user session: {:?}", session);
  225. match &session {
  226. None => {
  227. KV::remove(&self.config.session_cache_key).map_err(|e| UserError::new(ErrorCode::InternalError, &e))?
  228. },
  229. Some(session) => KV::set_str(&self.config.session_cache_key, session.clone().into()),
  230. }
  231. *self.session.write() = session;
  232. Ok(())
  233. }
  234. fn get_session(&self) -> Result<Session, UserError> {
  235. let mut session = { (*self.session.read()).clone() };
  236. if session.is_none() {
  237. match KV::get_str(&self.config.session_cache_key) {
  238. None => {},
  239. Some(s) => {
  240. session = Some(Session::from(s));
  241. let _ = self.set_session(session.clone())?;
  242. },
  243. }
  244. }
  245. match session {
  246. None => Err(UserError::unauthorized()),
  247. Some(session) => Ok(session),
  248. }
  249. }
  250. fn is_login(&self, email: &str) -> bool {
  251. match self.get_session() {
  252. Ok(session) => session.email == email,
  253. Err(_) => false,
  254. }
  255. }
  256. #[tracing::instrument(level = "debug", skip(self, token))]
  257. pub async fn start_ws_connection(&self, token: &str) -> Result<(), UserError> {
  258. if cfg!(feature = "http_server") {
  259. let addr = format!("{}/{}", self.server.ws_addr(), token);
  260. self.listen_on_websocket();
  261. let _ = self.ws_controller.start_connect(addr).await?;
  262. }
  263. Ok(())
  264. }
  265. #[tracing::instrument(level = "debug", skip(self))]
  266. fn listen_on_websocket(&self) {
  267. let mut notify = self.ws_controller.state_subscribe();
  268. let ws_controller = self.ws_controller.clone();
  269. let _ = tokio::spawn(async move {
  270. loop {
  271. match notify.recv().await {
  272. Ok(state) => {
  273. tracing::info!("Websocket state changed: {}", state);
  274. match state {
  275. WsState::Init => {},
  276. WsState::Connected(_) => {},
  277. WsState::Disconnected(_) => match ws_controller.retry().await {
  278. Ok(_) => {},
  279. Err(e) => {
  280. log::error!("websocket connect failed: {:?}", e);
  281. },
  282. },
  283. }
  284. },
  285. Err(e) => {
  286. log::error!("Websocket state notify error: {:?}", e);
  287. break;
  288. },
  289. }
  290. }
  291. });
  292. }
  293. }
  294. pub async fn update_user(
  295. _server: Server,
  296. pool: Arc<ConnectionPool>,
  297. params: UpdateUserParams,
  298. ) -> Result<(), UserError> {
  299. let changeset = UserTableChangeset::new(params);
  300. let conn = pool.get()?;
  301. diesel_update_table!(user_table, changeset, &*conn);
  302. Ok(())
  303. }
  304. impl UserDatabaseConnection for UserSession {
  305. fn get_connection(&self) -> Result<DBConnection, String> { self.db_connection().map_err(|e| format!("{:?}", e)) }
  306. }
  307. #[derive(Debug, Clone, Default, Serialize, Deserialize)]
  308. struct Session {
  309. user_id: String,
  310. token: String,
  311. email: String,
  312. }
  313. impl Session {
  314. pub fn new(user_id: &str, token: &str, email: &str) -> Self {
  315. Self {
  316. user_id: user_id.to_owned(),
  317. token: token.to_owned(),
  318. email: email.to_owned(),
  319. }
  320. }
  321. pub fn into_part(self) -> (String, String) { (self.user_id, self.token) }
  322. }
  323. impl std::convert::From<String> for Session {
  324. fn from(s: String) -> Self {
  325. match serde_json::from_str(&s) {
  326. Ok(s) => s,
  327. Err(e) => {
  328. log::error!("Deserialize string to Session failed: {:?}", e);
  329. Session::default()
  330. },
  331. }
  332. }
  333. }
  334. impl std::convert::From<Session> for String {
  335. fn from(session: Session) -> Self {
  336. match serde_json::to_string(&session) {
  337. Ok(s) => s,
  338. Err(e) => {
  339. log::error!("Serialize session to string failed: {:?}", e);
  340. "".to_string()
  341. },
  342. }
  343. }
  344. }