TpdsUpdateSender.js 6.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237
  1. const { ObjectId } = require('mongodb')
  2. const _ = require('lodash')
  3. const { callbackify } = require('util')
  4. const logger = require('logger-sharelatex')
  5. const metrics = require('@overleaf/metrics')
  6. const path = require('path')
  7. const request = require('request-promise-native')
  8. const settings = require('@overleaf/settings')
  9. const CollaboratorsGetter = require('../Collaborators/CollaboratorsGetter')
  10. .promises
  11. const UserGetter = require('../User/UserGetter.js').promises
  12. const tpdsUrl = _.get(settings, ['apis', 'thirdPartyDataStore', 'url'])
  13. async function addDoc(options) {
  14. metrics.inc('tpds.add-doc')
  15. options.streamOrigin =
  16. settings.apis.docstore.pubUrl +
  17. path.join(
  18. `/project/${options.project_id}`,
  19. `/doc/${options.doc_id}`,
  20. '/raw'
  21. )
  22. return addEntity(options)
  23. }
  24. async function addEntity(options) {
  25. const projectUserIds = await getProjectUsersIds(options.project_id)
  26. for (const userId of projectUserIds) {
  27. const job = {
  28. method: 'post',
  29. headers: {
  30. sl_entity_rev: options.rev,
  31. sl_project_id: options.project_id,
  32. sl_all_user_ids: JSON.stringify([userId]),
  33. sl_project_owner_user_id: projectUserIds[0],
  34. },
  35. uri: buildTpdsUrl(userId, options.project_name, options.path),
  36. title: 'addFile',
  37. streamOrigin: options.streamOrigin,
  38. }
  39. await enqueue(userId, 'pipeStreamFrom', job)
  40. }
  41. }
  42. async function addFile(options) {
  43. metrics.inc('tpds.add-file')
  44. options.streamOrigin =
  45. settings.apis.filestore.url +
  46. path.join(`/project/${options.project_id}`, `/file/${options.file_id}`)
  47. return addEntity(options)
  48. }
  49. function buildMovePaths(options) {
  50. if (options.newProjectName) {
  51. return {
  52. startPath: path.join('/', options.project_name, '/'),
  53. endPath: path.join('/', options.newProjectName, '/'),
  54. }
  55. } else {
  56. return {
  57. startPath: path.join('/', options.project_name, '/', options.startPath),
  58. endPath: path.join('/', options.project_name, '/', options.endPath),
  59. }
  60. }
  61. }
  62. function buildTpdsUrl(userId, projectName, filePath) {
  63. const projectPath = encodeURIComponent(path.join(projectName, '/', filePath))
  64. return `${tpdsUrl}/user/${userId}/entity/${projectPath}`
  65. }
  66. async function deleteEntity(options) {
  67. metrics.inc('tpds.delete-entity')
  68. const projectUserIds = await getProjectUsersIds(options.project_id)
  69. for (const userId of projectUserIds) {
  70. const job = {
  71. method: 'delete',
  72. headers: {
  73. sl_project_id: options.project_id,
  74. sl_all_user_ids: JSON.stringify([userId]),
  75. sl_project_owner_user_id: projectUserIds[0],
  76. },
  77. uri: buildTpdsUrl(userId, options.project_name, options.path),
  78. title: 'deleteEntity',
  79. sl_all_user_ids: JSON.stringify([userId]),
  80. }
  81. await enqueue(userId, 'standardHttpRequest', job)
  82. }
  83. }
  84. async function deleteProject(options) {
  85. // deletion only applies to project archiver
  86. const projectArchiverUrl = _.get(settings, [
  87. 'apis',
  88. 'project_archiver',
  89. 'url',
  90. ])
  91. // silently do nothing if project archiver url is not in settings
  92. if (!projectArchiverUrl) {
  93. return
  94. }
  95. metrics.inc('tpds.delete-project')
  96. // send the request directly to project archiver, bypassing third-party-datastore
  97. try {
  98. const response = await request({
  99. uri: `${settings.apis.project_archiver.url}/project/${options.project_id}`,
  100. method: 'delete',
  101. })
  102. return response
  103. } catch (err) {
  104. logger.error(
  105. { err, project_id: options.project_id },
  106. 'error deleting project in third party datastore (project_archiver)'
  107. )
  108. }
  109. }
  110. async function enqueue(group, method, job) {
  111. const tpdsWorkerUrl = _.get(settings, ['apis', 'tpdsworker', 'url'])
  112. // silently do nothing if worker url is not in settings
  113. if (!tpdsWorkerUrl) {
  114. return
  115. }
  116. try {
  117. const response = await request({
  118. uri: `${tpdsWorkerUrl}/enqueue/web_to_tpds_http_requests`,
  119. json: { group, job, method },
  120. method: 'post',
  121. timeout: 5 * 1000,
  122. })
  123. return response
  124. } catch (err) {
  125. // log error and continue
  126. logger.error({ err, group, job, method }, 'error enqueueing tpdsworker job')
  127. }
  128. }
  129. async function getProjectUsersIds(projectId) {
  130. // get list of all user ids with access to project. project owner
  131. // will always be the first entry in the list.
  132. const [
  133. ownerUserId,
  134. ...invitedUserIds
  135. ] = await CollaboratorsGetter.getInvitedMemberIds(projectId)
  136. // if there are no invited users, always return the owner
  137. if (!invitedUserIds.length) {
  138. return [ownerUserId]
  139. }
  140. // filter invited users to only return those with dropbox linked
  141. const dropboxUsers = await UserGetter.getUsers(
  142. {
  143. _id: { $in: invitedUserIds.map(id => ObjectId(id)) },
  144. 'dropbox.access_token.uid': { $ne: null },
  145. },
  146. {
  147. _id: 1,
  148. }
  149. )
  150. const dropboxUserIds = dropboxUsers.map(user => user._id)
  151. return [ownerUserId, ...dropboxUserIds]
  152. }
  153. async function moveEntity(options) {
  154. metrics.inc('tpds.move-entity')
  155. const projectUserIds = await getProjectUsersIds(options.project_id)
  156. const { endPath, startPath } = buildMovePaths(options)
  157. for (const userId of projectUserIds) {
  158. const job = {
  159. method: 'put',
  160. title: 'moveEntity',
  161. uri: `${tpdsUrl}/user/${userId}/entity`,
  162. headers: {
  163. sl_project_id: options.project_id,
  164. sl_entity_rev: options.rev,
  165. sl_all_user_ids: JSON.stringify([userId]),
  166. sl_project_owner_user_id: projectUserIds[0],
  167. },
  168. json: {
  169. user_id: userId,
  170. endPath,
  171. startPath,
  172. },
  173. }
  174. await enqueue(userId, 'standardHttpRequest', job)
  175. }
  176. }
  177. async function pollDropboxForUser(userId) {
  178. metrics.inc('tpds.poll-dropbox')
  179. const job = {
  180. method: 'post',
  181. uri: `${tpdsUrl}/user/poll`,
  182. json: {
  183. user_ids: [userId],
  184. },
  185. }
  186. return enqueue(`poll-dropbox:${userId}`, 'standardHttpRequest', job)
  187. }
  188. const TpdsUpdateSender = {
  189. addDoc: callbackify(addDoc),
  190. addEntity: callbackify(addEntity),
  191. addFile: callbackify(addFile),
  192. deleteEntity: callbackify(deleteEntity),
  193. deleteProject: callbackify(deleteProject),
  194. enqueue: callbackify(enqueue),
  195. moveEntity: callbackify(moveEntity),
  196. pollDropboxForUser: callbackify(pollDropboxForUser),
  197. promises: {
  198. addDoc,
  199. addEntity,
  200. addFile,
  201. deleteEntity,
  202. deleteProject,
  203. enqueue,
  204. moveEntity,
  205. pollDropboxForUser,
  206. },
  207. }
  208. module.exports = TpdsUpdateSender