batchedUpdate.js 7.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321
  1. // @ts-check
  2. /* eslint-disable no-console */
  3. const { ObjectId, ReadPreference } = require('mongodb')
  4. const READ_PREFERENCE_SECONDARY =
  5. process.env.MONGO_HAS_SECONDARIES === 'true'
  6. ? ReadPreference.secondary.mode
  7. : ReadPreference.secondaryPreferred.mode
  8. const ONE_MONTH_IN_MS = 1000 * 60 * 60 * 24 * 31
  9. let ID_EDGE_PAST
  10. const ID_EDGE_FUTURE = objectIdFromMs(Date.now() + 1000)
  11. let BATCH_DESCENDING
  12. let BATCH_SIZE
  13. let VERBOSE_LOGGING
  14. let BATCH_RANGE_START
  15. let BATCH_RANGE_END
  16. let BATCH_MAX_TIME_SPAN_IN_MS
  17. let BATCHED_UPDATE_RUNNING = false
  18. /**
  19. * @typedef {import("mongodb").Collection} Collection
  20. * @typedef {import("mongodb-legacy").Collection} LegacyCollection
  21. * @typedef {import("mongodb").Document} Document
  22. * @typedef {import("mongodb").FindOptions} FindOptions
  23. * @typedef {import("mongodb").UpdateFilter<Document>} UpdateDocument
  24. */
  25. /**
  26. * @typedef {Object} BatchedUpdateOptions
  27. * @property {string} [BATCH_DESCENDING]
  28. * @property {string} [BATCH_LAST_ID]
  29. * @property {string} [BATCH_MAX_TIME_SPAN_IN_MS]
  30. * @property {string} [BATCH_RANGE_END]
  31. * @property {string} [BATCH_RANGE_START]
  32. * @property {string} [BATCH_SIZE]
  33. * @property {string} [VERBOSE_LOGGING]
  34. * @property {(progress: string) => Promise<void>} [trackProgress]
  35. */
  36. /**
  37. * @param {BatchedUpdateOptions} options
  38. */
  39. function refreshGlobalOptionsForBatchedUpdate(options = {}) {
  40. options = Object.assign({}, options, process.env)
  41. BATCH_DESCENDING = options.BATCH_DESCENDING === 'true'
  42. BATCH_SIZE = parseInt(options.BATCH_SIZE || '1000', 10) || 1000
  43. VERBOSE_LOGGING = options.VERBOSE_LOGGING === 'true'
  44. if (options.BATCH_LAST_ID) {
  45. BATCH_RANGE_START = objectIdFromInput(options.BATCH_LAST_ID)
  46. } else if (options.BATCH_RANGE_START) {
  47. BATCH_RANGE_START = objectIdFromInput(options.BATCH_RANGE_START)
  48. } else {
  49. if (BATCH_DESCENDING) {
  50. BATCH_RANGE_START = ID_EDGE_FUTURE
  51. } else {
  52. BATCH_RANGE_START = ID_EDGE_PAST
  53. }
  54. }
  55. BATCH_MAX_TIME_SPAN_IN_MS = parseInt(
  56. options.BATCH_MAX_TIME_SPAN_IN_MS || ONE_MONTH_IN_MS.toString(),
  57. 10
  58. )
  59. if (options.BATCH_RANGE_END) {
  60. BATCH_RANGE_END = objectIdFromInput(options.BATCH_RANGE_END)
  61. } else {
  62. if (BATCH_DESCENDING) {
  63. BATCH_RANGE_END = ID_EDGE_PAST
  64. } else {
  65. BATCH_RANGE_END = ID_EDGE_FUTURE
  66. }
  67. }
  68. }
  69. /**
  70. * @param {Collection | LegacyCollection} collection
  71. * @param {Document} query
  72. * @param {ObjectId} start
  73. * @param {ObjectId} end
  74. * @param {Document} projection
  75. * @param {FindOptions} findOptions
  76. * @return {Promise<Array<Document>>}
  77. */
  78. async function getNextBatch(
  79. collection,
  80. query,
  81. start,
  82. end,
  83. projection,
  84. findOptions
  85. ) {
  86. if (BATCH_DESCENDING) {
  87. query._id = {
  88. $gt: end,
  89. $lte: start,
  90. }
  91. } else {
  92. query._id = {
  93. $gt: start,
  94. $lte: end,
  95. }
  96. }
  97. return await collection
  98. .find(query, findOptions)
  99. .project(projection)
  100. .sort({ _id: BATCH_DESCENDING ? -1 : 1 })
  101. .limit(BATCH_SIZE)
  102. .toArray()
  103. }
  104. /**
  105. * @param {Collection | LegacyCollection} collection
  106. * @param {Array<Document>} nextBatch
  107. * @param {UpdateDocument} update
  108. * @return {Promise<void>}
  109. */
  110. async function performUpdate(collection, nextBatch, update) {
  111. await collection.updateMany(
  112. { _id: { $in: nextBatch.map(entry => entry._id) } },
  113. update
  114. )
  115. }
  116. /**
  117. * @param {string} input
  118. * @return {ObjectId}
  119. */
  120. function objectIdFromInput(input) {
  121. if (input.includes('T')) {
  122. const t = new Date(input).getTime()
  123. if (Number.isNaN(t)) throw new Error(`${input} is not a valid date`)
  124. return objectIdFromMs(t)
  125. } else {
  126. return new ObjectId(input)
  127. }
  128. }
  129. /**
  130. * @param {ObjectId} objectId
  131. * @return {string}
  132. */
  133. function renderObjectId(objectId) {
  134. return `${objectId} (${objectId.getTimestamp().toISOString()})`
  135. }
  136. /**
  137. * @param {number} ms
  138. * @return {ObjectId}
  139. */
  140. function objectIdFromMs(ms) {
  141. return ObjectId.createFromTime(ms / 1000)
  142. }
  143. /**
  144. * @param {ObjectId} id
  145. * @return {number}
  146. */
  147. function getMsFromObjectId(id) {
  148. return id.getTimestamp().getTime()
  149. }
  150. /**
  151. * @param {ObjectId} start
  152. * @return {ObjectId}
  153. */
  154. function getNextEnd(start) {
  155. let end
  156. if (BATCH_DESCENDING) {
  157. end = objectIdFromMs(getMsFromObjectId(start) - BATCH_MAX_TIME_SPAN_IN_MS)
  158. if (getMsFromObjectId(end) <= getMsFromObjectId(BATCH_RANGE_END)) {
  159. end = BATCH_RANGE_END
  160. }
  161. } else {
  162. end = objectIdFromMs(getMsFromObjectId(start) + BATCH_MAX_TIME_SPAN_IN_MS)
  163. if (getMsFromObjectId(end) >= getMsFromObjectId(BATCH_RANGE_END)) {
  164. end = BATCH_RANGE_END
  165. }
  166. }
  167. return end
  168. }
  169. /**
  170. * @param {Collection | LegacyCollection} collection
  171. * @return {Promise<ObjectId|null>}
  172. */
  173. async function getIdEdgePast(collection) {
  174. const [first] = await collection
  175. .find({})
  176. .project({ _id: 1 })
  177. .sort({ _id: 1 })
  178. .limit(1)
  179. .toArray()
  180. if (!first) return null
  181. // Go one second further into the past in order to include the first entry via
  182. // first._id > ID_EDGE_PAST
  183. return objectIdFromMs(Math.max(0, getMsFromObjectId(first._id) - 1000))
  184. }
  185. /**
  186. * @param {Collection | LegacyCollection} collection
  187. * @param {Document} query
  188. * @param {UpdateDocument | ((batch: Array<Document>) => Promise<void>)} update
  189. * @param {Document} [projection]
  190. * @param {FindOptions} [findOptions]
  191. * @param {BatchedUpdateOptions} [batchedUpdateOptions]
  192. */
  193. async function batchedUpdate(
  194. collection,
  195. query,
  196. update,
  197. projection,
  198. findOptions,
  199. batchedUpdateOptions = {}
  200. ) {
  201. // only a single batchedUpdate can run at a time due to global variables
  202. if (BATCHED_UPDATE_RUNNING) {
  203. throw new Error('batchedUpdate is already running')
  204. }
  205. try {
  206. BATCHED_UPDATE_RUNNING = true
  207. ID_EDGE_PAST = await getIdEdgePast(collection)
  208. if (!ID_EDGE_PAST) {
  209. console.warn(
  210. `The collection ${collection.collectionName} appears to be empty.`
  211. )
  212. return 0
  213. }
  214. refreshGlobalOptionsForBatchedUpdate(batchedUpdateOptions)
  215. const { trackProgress = async progress => console.warn(progress) } =
  216. batchedUpdateOptions
  217. findOptions = findOptions || {}
  218. findOptions.readPreference = READ_PREFERENCE_SECONDARY
  219. projection = projection || { _id: 1 }
  220. let nextBatch
  221. let updated = 0
  222. let start = BATCH_RANGE_START
  223. while (start !== BATCH_RANGE_END) {
  224. let end = getNextEnd(start)
  225. nextBatch = await getNextBatch(
  226. collection,
  227. query,
  228. start,
  229. end,
  230. projection,
  231. findOptions
  232. )
  233. if (nextBatch.length > 0) {
  234. end = nextBatch[nextBatch.length - 1]._id
  235. updated += nextBatch.length
  236. if (VERBOSE_LOGGING) {
  237. console.log(
  238. `Running update on batch with ids ${JSON.stringify(
  239. nextBatch.map(entry => entry._id)
  240. )}`
  241. )
  242. }
  243. await trackProgress(
  244. `Running update on batch ending ${renderObjectId(end)}`
  245. )
  246. if (typeof update === 'function') {
  247. await update(nextBatch)
  248. } else {
  249. await performUpdate(collection, nextBatch, update)
  250. }
  251. }
  252. await trackProgress(`Completed batch ending ${renderObjectId(end)}`)
  253. start = end
  254. }
  255. return updated
  256. } finally {
  257. BATCHED_UPDATE_RUNNING = false
  258. }
  259. }
  260. /**
  261. * @param {Collection | LegacyCollection} collection
  262. * @param {Document} query
  263. * @param {UpdateDocument | ((batch: Array<Object>) => Promise<void>)} update
  264. * @param {Document} [projection]
  265. * @param {FindOptions} [findOptions]
  266. * @param {BatchedUpdateOptions} [batchedUpdateOptions]
  267. */
  268. function batchedUpdateWithResultHandling(
  269. collection,
  270. query,
  271. update,
  272. projection,
  273. findOptions,
  274. batchedUpdateOptions
  275. ) {
  276. batchedUpdate(
  277. collection,
  278. query,
  279. update,
  280. projection,
  281. findOptions,
  282. batchedUpdateOptions
  283. )
  284. .then(processed => {
  285. console.error({ processed })
  286. process.exit(0)
  287. })
  288. .catch(error => {
  289. console.error({ error })
  290. process.exit(1)
  291. })
  292. }
  293. module.exports = {
  294. READ_PREFERENCE_SECONDARY,
  295. objectIdFromInput,
  296. renderObjectId,
  297. batchedUpdate,
  298. batchedUpdateWithResultHandling,
  299. }