delete_orphaned_docs_online_check.js 5.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181
  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('../app/src/util/promises')
  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 = 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(ObjectId)
  63. console.log('Checking projects', JSON.stringify(projectIds))
  64. const { nProjectsWithOrphanedDocs, nDeletedDocs } = await processBatch(
  65. projectIds
  66. )
  67. nProjectsProcessedTotal += projectIds.length
  68. nProjectsWithOrphanedDocsTotal += nProjectsWithOrphanedDocs
  69. nDeletedDocsTotal += nDeletedDocs
  70. if (docs.length === BATCH_SIZE) {
  71. // This project may have more than BATCH_SIZE docs.
  72. const lastDoc = docs[docs.length - 1]
  73. // Resume from after this projectId.
  74. upperProjectId = lastDoc.project_id
  75. }
  76. }
  77. console.error(
  78. 'Processed %d projects ' +
  79. '(%d projects with orphaned docs/%d docs deleted) ' +
  80. 'until %s',
  81. nProjectsProcessedTotal,
  82. nProjectsWithOrphanedDocsTotal,
  83. nDeletedDocsTotal,
  84. upperProjectId
  85. )
  86. lowerProjectId = upperProjectId
  87. }
  88. }
  89. async function getProjectDocs(projectId) {
  90. return await db.docs
  91. .find(
  92. { project_id: projectId },
  93. {
  94. projection: { _id: 1 },
  95. readPreference: READ_PREFERENCE_PRIMARY,
  96. }
  97. )
  98. .toArray()
  99. }
  100. async function processBatch(projectIds) {
  101. const projectsWithOrphanedDocs = await getHardDeletedProjectIds({
  102. projectIds,
  103. READ_CONCURRENCY_PRIMARY,
  104. READ_CONCURRENCY_SECONDARY,
  105. })
  106. let nDeletedDocs = 0
  107. async function countOrphanedDocs(projectId) {
  108. const docs = await getProjectDocs(projectId)
  109. nDeletedDocs += docs.length
  110. console.log(
  111. 'Deleted project %s has %s orphaned docs: %s',
  112. projectId,
  113. docs.length,
  114. JSON.stringify(docs.map(doc => doc._id))
  115. )
  116. }
  117. await promiseMapWithLimit(
  118. READ_CONCURRENCY_PRIMARY,
  119. projectsWithOrphanedDocs,
  120. countOrphanedDocs
  121. )
  122. if (!DRY_RUN) {
  123. await promiseMapWithLimit(
  124. WRITE_CONCURRENCY,
  125. projectsWithOrphanedDocs,
  126. DocstoreManager.promises.destroyProject
  127. )
  128. }
  129. const nProjectsWithOrphanedDocs = projectsWithOrphanedDocs.length
  130. return { nProjectsWithOrphanedDocs, nDeletedDocs }
  131. }
  132. async function letUserDoubleCheckInputs() {
  133. console.error(
  134. 'Options:',
  135. JSON.stringify(
  136. {
  137. BATCH_LAST_ID,
  138. BATCH_SIZE,
  139. DRY_RUN,
  140. INCREMENT_BY_S,
  141. STOP_AT_S,
  142. READ_CONCURRENCY_SECONDARY,
  143. READ_CONCURRENCY_PRIMARY,
  144. WRITE_CONCURRENCY,
  145. LET_USER_DOUBLE_CHECK_INPUTS_FOR,
  146. },
  147. null,
  148. 2
  149. )
  150. )
  151. console.error(
  152. 'Waiting for you to double check inputs for',
  153. LET_USER_DOUBLE_CHECK_INPUTS_FOR,
  154. 'ms'
  155. )
  156. await sleep(LET_USER_DOUBLE_CHECK_INPUTS_FOR)
  157. }
  158. main()
  159. .then(() => {
  160. console.error('Done.')
  161. process.exit(0)
  162. })
  163. .catch(error => {
  164. console.error({ error })
  165. process.exit(1)
  166. })