delete_orphaned_docs_online_check.js 7.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257
  1. const DocstoreManager = require('../app/src/Features/Docstore/DocstoreManager')
  2. const { promisify } = require('util')
  3. const { ObjectId, ReadPreference } = require('mongodb')
  4. const { db, waitForDb } = require('../app/src/infrastructure/mongodb')
  5. const { promiseMapWithLimit } = require('../app/src/util/promises')
  6. const sleep = promisify(setTimeout)
  7. const NOW_IN_S = Date.now() / 1000
  8. const ONE_WEEK_IN_S = 60 * 60 * 24 * 7
  9. const TEN_SECONDS = 10 * 1000
  10. const DRY_RUN = process.env.DRY_RUN === 'true'
  11. if (!process.env.BATCH_LAST_ID) {
  12. console.error('Set BATCH_LAST_ID and re-run.')
  13. process.exit(1)
  14. }
  15. const BATCH_LAST_ID = ObjectId(process.env.BATCH_LAST_ID)
  16. const INCREMENT_BY_S = parseInt(process.env.INCREMENT_BY_S, 10) || ONE_WEEK_IN_S
  17. const BATCH_SIZE = parseInt(process.env.BATCH_SIZE, 10) || 1000
  18. const READ_CONCURRENCY_SECONDARY =
  19. parseInt(process.env.READ_CONCURRENCY_SECONDARY, 10) || 1000
  20. const READ_CONCURRENCY_PRIMARY =
  21. parseInt(process.env.READ_CONCURRENCY_PRIMARY, 10) || 500
  22. const STOP_AT_S = parseInt(process.env.STOP_AT_S, 10) || NOW_IN_S
  23. const WRITE_CONCURRENCY = parseInt(process.env.WRITE_CONCURRENCY, 10) || 10
  24. const LET_USER_DOUBLE_CHECK_INPUTS_FOR =
  25. parseInt(process.env.LET_USER_DOUBLE_CHECK_INPUTS_FOR, 10) || TEN_SECONDS
  26. function getSecondsFromObjectId(id) {
  27. return id.getTimestamp().getTime() / 1000
  28. }
  29. async function main() {
  30. await letUserDoubleCheckInputs()
  31. await waitForDb()
  32. let lowerProjectId = BATCH_LAST_ID
  33. let nProjectsProcessedTotal = 0
  34. let nProjectsWithOrphanedDocsTotal = 0
  35. let nDeletedDocsTotal = 0
  36. while (getSecondsFromObjectId(lowerProjectId) <= STOP_AT_S) {
  37. const upperTime = getSecondsFromObjectId(lowerProjectId) + INCREMENT_BY_S
  38. let upperProjectId = ObjectId.createFromTime(upperTime)
  39. const query = {
  40. project_id: {
  41. // exclude edge
  42. $gt: lowerProjectId,
  43. // include edge
  44. $lte: upperProjectId,
  45. },
  46. }
  47. const docs = await db.docs
  48. .find(query, { readPreference: ReadPreference.SECONDARY })
  49. .project({ project_id: 1 })
  50. .sort({ project_id: 1 })
  51. .limit(BATCH_SIZE)
  52. .toArray()
  53. if (docs.length) {
  54. const projectIds = Array.from(
  55. new Set(docs.map(doc => doc.project_id.toString()))
  56. ).map(ObjectId)
  57. console.log('Checking projects', JSON.stringify(projectIds))
  58. const { nProjectsWithOrphanedDocs, nDeletedDocs } = await processBatch(
  59. projectIds
  60. )
  61. nProjectsProcessedTotal += projectIds.length
  62. nProjectsWithOrphanedDocsTotal += nProjectsWithOrphanedDocs
  63. nDeletedDocsTotal += nDeletedDocs
  64. if (docs.length === BATCH_SIZE) {
  65. // This project may have more than BATCH_SIZE docs.
  66. const lastDoc = docs[docs.length - 1]
  67. // Resume from after this projectId.
  68. upperProjectId = lastDoc.project_id
  69. }
  70. }
  71. console.error(
  72. 'Processed %d projects ' +
  73. '(%d projects with orphaned docs/%d docs deleted) ' +
  74. 'until %s',
  75. nProjectsProcessedTotal,
  76. nProjectsWithOrphanedDocsTotal,
  77. nDeletedDocsTotal,
  78. upperProjectId
  79. )
  80. lowerProjectId = upperProjectId
  81. }
  82. }
  83. async function getDeletedProject(projectId, readPreference) {
  84. return await db.deletedProjects.findOne(
  85. { 'deleterData.deletedProjectId': projectId },
  86. {
  87. // There is no index on .project. Pull down something small.
  88. projection: { 'project._id': 1 },
  89. readPreference,
  90. }
  91. )
  92. }
  93. async function getProject(projectId, readPreference) {
  94. return await db.projects.findOne(
  95. { _id: projectId },
  96. {
  97. // Pulling down an empty object is fine for differentiating with null.
  98. projection: { _id: 0 },
  99. readPreference,
  100. }
  101. )
  102. }
  103. async function getProjectDocs(projectId) {
  104. return await db.docs
  105. .find(
  106. { project_id: projectId },
  107. {
  108. projection: { _id: 1 },
  109. readPreference: ReadPreference.PRIMARY,
  110. }
  111. )
  112. .toArray()
  113. }
  114. async function checkProjectExistsWithReadPreference(projectId, readPreference) {
  115. // NOTE: Possible race conditions!
  116. // There are two processes which are racing with our queries:
  117. // 1. project deletion
  118. // 2. project restoring
  119. // For 1. we check the projects collection before deletedProjects.
  120. // If a project were to be delete in this very moment, we should see the
  121. // soft-deleted entry which is created before deleting the projects entry.
  122. // For 2. we check the projects collection after deletedProjects again.
  123. // If a project were to be restored in this very moment, it is very likely
  124. // to see the projects entry again.
  125. // Unlikely edge case: Restore+Deletion in rapid succession.
  126. // We could add locking to the ProjectDeleter for ruling ^ out.
  127. if (await getProject(projectId, readPreference)) {
  128. // The project is live.
  129. return true
  130. }
  131. const deletedProject = await getDeletedProject(projectId, readPreference)
  132. if (deletedProject && deletedProject.project) {
  133. // The project is registered for hard-deletion.
  134. return true
  135. }
  136. if (await getProject(projectId, readPreference)) {
  137. // The project was just restored.
  138. return true
  139. }
  140. // The project does not exist.
  141. return false
  142. }
  143. async function checkProjectExistsOnPrimary(projectId) {
  144. return await checkProjectExistsWithReadPreference(
  145. projectId,
  146. ReadPreference.PRIMARY
  147. )
  148. }
  149. async function checkProjectExistsOnSecondary(projectId) {
  150. return await checkProjectExistsWithReadPreference(
  151. projectId,
  152. ReadPreference.SECONDARY
  153. )
  154. }
  155. async function processBatch(projectIds) {
  156. const doubleCheckProjectIdsOnPrimary = []
  157. let nDeletedDocs = 0
  158. async function checkProjectOnSecondary(projectId) {
  159. if (await checkProjectExistsOnSecondary(projectId)) {
  160. // Finding a project with secondary confidence is sufficient.
  161. return
  162. }
  163. // At this point, the secondaries deem this project as having orphaned docs.
  164. doubleCheckProjectIdsOnPrimary.push(projectId)
  165. }
  166. const projectsWithOrphanedDocs = []
  167. async function checkProjectOnPrimary(projectId) {
  168. if (await checkProjectExistsOnPrimary(projectId)) {
  169. // The project is actually live.
  170. return
  171. }
  172. projectsWithOrphanedDocs.push(projectId)
  173. const docs = await getProjectDocs(projectId)
  174. nDeletedDocs += docs.length
  175. console.log(
  176. 'Deleted project %s has %s orphaned docs: %s',
  177. projectId,
  178. docs.length,
  179. JSON.stringify(docs.map(doc => doc._id))
  180. )
  181. }
  182. await promiseMapWithLimit(
  183. READ_CONCURRENCY_SECONDARY,
  184. projectIds,
  185. checkProjectOnSecondary
  186. )
  187. await promiseMapWithLimit(
  188. READ_CONCURRENCY_PRIMARY,
  189. doubleCheckProjectIdsOnPrimary,
  190. checkProjectOnPrimary
  191. )
  192. if (!DRY_RUN) {
  193. await promiseMapWithLimit(
  194. WRITE_CONCURRENCY,
  195. projectsWithOrphanedDocs,
  196. DocstoreManager.promises.destroyProject
  197. )
  198. }
  199. const nProjectsWithOrphanedDocs = projectsWithOrphanedDocs.length
  200. return { nProjectsWithOrphanedDocs, nDeletedDocs }
  201. }
  202. async function letUserDoubleCheckInputs() {
  203. console.error(
  204. 'Options:',
  205. JSON.stringify(
  206. {
  207. BATCH_LAST_ID,
  208. BATCH_SIZE,
  209. DRY_RUN,
  210. INCREMENT_BY_S,
  211. STOP_AT_S,
  212. READ_CONCURRENCY_SECONDARY,
  213. READ_CONCURRENCY_PRIMARY,
  214. WRITE_CONCURRENCY,
  215. LET_USER_DOUBLE_CHECK_INPUTS_FOR,
  216. },
  217. null,
  218. 2
  219. )
  220. )
  221. console.error(
  222. 'Waiting for you to double check inputs for',
  223. LET_USER_DOUBLE_CHECK_INPUTS_FOR,
  224. 'ms'
  225. )
  226. await sleep(LET_USER_DOUBLE_CHECK_INPUTS_FOR)
  227. }
  228. main()
  229. .then(() => {
  230. console.error('Done.')
  231. process.exit(0)
  232. })
  233. .catch(error => {
  234. console.error({ error })
  235. process.exit(1)
  236. })