delete_orphaned_docs_online_check.js 5.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180
  1. const DocstoreManager = require('../app/src/Features/Docstore/DocstoreManager')
  2. const { promisify } = require('util')
  3. const { ObjectId } = require('mongodb')
  4. const {
  5. db,
  6. waitForDb,
  7. READ_PREFERENCE_PRIMARY,
  8. READ_PREFERENCE_SECONDARY,
  9. } = require('../app/src/infrastructure/mongodb')
  10. const { promiseMapWithLimit } = require('@overleaf/promise-utils')
  11. const { getHardDeletedProjectIds } = require('./delete_orphaned_data_helper')
  12. const sleep = promisify(setTimeout)
  13. const NOW_IN_S = Date.now() / 1000
  14. const ONE_WEEK_IN_S = 60 * 60 * 24 * 7
  15. const TEN_SECONDS = 10 * 1000
  16. const DRY_RUN = process.env.DRY_RUN === 'true'
  17. if (!process.env.BATCH_LAST_ID) {
  18. console.error('Set BATCH_LAST_ID and re-run.')
  19. process.exit(1)
  20. }
  21. const BATCH_LAST_ID = new ObjectId(process.env.BATCH_LAST_ID)
  22. const INCREMENT_BY_S = parseInt(process.env.INCREMENT_BY_S, 10) || ONE_WEEK_IN_S
  23. const BATCH_SIZE = parseInt(process.env.BATCH_SIZE, 10) || 1000
  24. const READ_CONCURRENCY_SECONDARY =
  25. parseInt(process.env.READ_CONCURRENCY_SECONDARY, 10) || 1000
  26. const READ_CONCURRENCY_PRIMARY =
  27. parseInt(process.env.READ_CONCURRENCY_PRIMARY, 10) || 500
  28. const STOP_AT_S = parseInt(process.env.STOP_AT_S, 10) || NOW_IN_S
  29. const WRITE_CONCURRENCY = parseInt(process.env.WRITE_CONCURRENCY, 10) || 10
  30. const LET_USER_DOUBLE_CHECK_INPUTS_FOR =
  31. parseInt(process.env.LET_USER_DOUBLE_CHECK_INPUTS_FOR, 10) || TEN_SECONDS
  32. function getSecondsFromObjectId(id) {
  33. return id.getTimestamp().getTime() / 1000
  34. }
  35. async function main() {
  36. await letUserDoubleCheckInputs()
  37. await waitForDb()
  38. let lowerProjectId = BATCH_LAST_ID
  39. let nProjectsProcessedTotal = 0
  40. let nProjectsWithOrphanedDocsTotal = 0
  41. let nDeletedDocsTotal = 0
  42. while (getSecondsFromObjectId(lowerProjectId) <= STOP_AT_S) {
  43. const upperTime = getSecondsFromObjectId(lowerProjectId) + INCREMENT_BY_S
  44. let upperProjectId = ObjectId.createFromTime(upperTime)
  45. const query = {
  46. project_id: {
  47. // exclude edge
  48. $gt: lowerProjectId,
  49. // include edge
  50. $lte: upperProjectId,
  51. },
  52. }
  53. const docs = await db.docs
  54. .find(query, { readPreference: READ_PREFERENCE_SECONDARY })
  55. .project({ project_id: 1 })
  56. .sort({ project_id: 1 })
  57. .limit(BATCH_SIZE)
  58. .toArray()
  59. if (docs.length) {
  60. const projectIds = Array.from(
  61. new Set(docs.map(doc => doc.project_id.toString()))
  62. ).map(id => new ObjectId(id))
  63. console.log('Checking projects', JSON.stringify(projectIds))
  64. const { nProjectsWithOrphanedDocs, nDeletedDocs } =
  65. await processBatch(projectIds)
  66. nProjectsProcessedTotal += projectIds.length
  67. nProjectsWithOrphanedDocsTotal += nProjectsWithOrphanedDocs
  68. nDeletedDocsTotal += nDeletedDocs
  69. if (docs.length === BATCH_SIZE) {
  70. // This project may have more than BATCH_SIZE docs.
  71. const lastDoc = docs[docs.length - 1]
  72. // Resume from after this projectId.
  73. upperProjectId = lastDoc.project_id
  74. }
  75. }
  76. console.error(
  77. 'Processed %d projects ' +
  78. '(%d projects with orphaned docs/%d docs deleted) ' +
  79. 'until %s',
  80. nProjectsProcessedTotal,
  81. nProjectsWithOrphanedDocsTotal,
  82. nDeletedDocsTotal,
  83. upperProjectId
  84. )
  85. lowerProjectId = upperProjectId
  86. }
  87. }
  88. async function getProjectDocs(projectId) {
  89. return await db.docs
  90. .find(
  91. { project_id: projectId },
  92. {
  93. projection: { _id: 1 },
  94. readPreference: READ_PREFERENCE_PRIMARY,
  95. }
  96. )
  97. .toArray()
  98. }
  99. async function processBatch(projectIds) {
  100. const projectsWithOrphanedDocs = await getHardDeletedProjectIds({
  101. projectIds,
  102. READ_CONCURRENCY_PRIMARY,
  103. READ_CONCURRENCY_SECONDARY,
  104. })
  105. let nDeletedDocs = 0
  106. async function countOrphanedDocs(projectId) {
  107. const docs = await getProjectDocs(projectId)
  108. nDeletedDocs += docs.length
  109. console.log(
  110. 'Deleted project %s has %s orphaned docs: %s',
  111. projectId,
  112. docs.length,
  113. JSON.stringify(docs.map(doc => doc._id))
  114. )
  115. }
  116. await promiseMapWithLimit(
  117. READ_CONCURRENCY_PRIMARY,
  118. projectsWithOrphanedDocs,
  119. countOrphanedDocs
  120. )
  121. if (!DRY_RUN) {
  122. await promiseMapWithLimit(
  123. WRITE_CONCURRENCY,
  124. projectsWithOrphanedDocs,
  125. DocstoreManager.promises.destroyProject
  126. )
  127. }
  128. const nProjectsWithOrphanedDocs = projectsWithOrphanedDocs.length
  129. return { nProjectsWithOrphanedDocs, nDeletedDocs }
  130. }
  131. async function letUserDoubleCheckInputs() {
  132. console.error(
  133. 'Options:',
  134. JSON.stringify(
  135. {
  136. BATCH_LAST_ID,
  137. BATCH_SIZE,
  138. DRY_RUN,
  139. INCREMENT_BY_S,
  140. STOP_AT_S,
  141. READ_CONCURRENCY_SECONDARY,
  142. READ_CONCURRENCY_PRIMARY,
  143. WRITE_CONCURRENCY,
  144. LET_USER_DOUBLE_CHECK_INPUTS_FOR,
  145. },
  146. null,
  147. 2
  148. )
  149. )
  150. console.error(
  151. 'Waiting for you to double check inputs for',
  152. LET_USER_DOUBLE_CHECK_INPUTS_FOR,
  153. 'ms'
  154. )
  155. await sleep(LET_USER_DOUBLE_CHECK_INPUTS_FOR)
  156. }
  157. main()
  158. .then(() => {
  159. console.error('Done.')
  160. process.exit(0)
  161. })
  162. .catch(error => {
  163. console.error({ error })
  164. process.exit(1)
  165. })