delete_orphaned_docs_online_check.js 5.0 KB

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