DispatchManager.js 3.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112
  1. /* eslint-disable
  2. no-unused-vars,
  3. */
  4. // TODO: This file was created by bulk-decaffeinate.
  5. // Fix any style issues and re-enable lint.
  6. /*
  7. * decaffeinate suggestions:
  8. * DS101: Remove unnecessary use of Array.from
  9. * DS102: Remove unnecessary code created because of implicit returns
  10. * DS202: Simplify dynamic range loops
  11. * DS205: Consider reworking code to avoid use of IIFEs
  12. * DS207: Consider shorter variations of null checks
  13. * Full docs: https://github.com/decaffeinate/decaffeinate/blob/master/docs/suggestions.md
  14. */
  15. let DispatchManager
  16. const Settings = require('@overleaf/settings')
  17. const logger = require('@overleaf/logger')
  18. const Keys = require('./UpdateKeys')
  19. const redis = require('@overleaf/redis-wrapper')
  20. const Errors = require('./Errors')
  21. const _ = require('lodash')
  22. const UpdateManager = require('./UpdateManager')
  23. const Metrics = require('./Metrics')
  24. const RateLimitManager = require('./RateLimitManager')
  25. module.exports = DispatchManager = {
  26. createDispatcher(RateLimiter, queueShardNumber) {
  27. let pendingListKey
  28. if (queueShardNumber === 0) {
  29. pendingListKey = 'pending-updates-list'
  30. } else {
  31. pendingListKey = `pending-updates-list-${queueShardNumber}`
  32. }
  33. const client = redis.createClient(Settings.redis.documentupdater)
  34. const worker = {
  35. client,
  36. _waitForUpdateThenDispatchWorker(callback) {
  37. if (callback == null) {
  38. callback = function () {}
  39. }
  40. const timer = new Metrics.Timer('worker.waiting')
  41. return worker.client.blpop(pendingListKey, 0, function (error, result) {
  42. logger.debug(`getting ${queueShardNumber}`, error, result)
  43. timer.done()
  44. if (error != null) {
  45. return callback(error)
  46. }
  47. if (result == null) {
  48. return callback()
  49. }
  50. const [listName, docKey] = Array.from(result)
  51. const [projectId, docId] = Array.from(
  52. Keys.splitProjectIdAndDocId(docKey)
  53. )
  54. // Dispatch this in the background
  55. const backgroundTask = cb =>
  56. UpdateManager.processOutstandingUpdatesWithLock(
  57. projectId,
  58. docId,
  59. function (error) {
  60. // log everything except OpRangeNotAvailable errors, these are normal
  61. if (error != null) {
  62. // downgrade OpRangeNotAvailable and "Delete component" errors so they are not sent to sentry
  63. const logAsDebug =
  64. error instanceof Errors.OpRangeNotAvailableError ||
  65. error instanceof Errors.DeleteMismatchError
  66. if (logAsDebug) {
  67. logger.debug(
  68. { err: error, projectId, docId },
  69. 'error processing update'
  70. )
  71. } else {
  72. logger.error(
  73. { err: error, projectId, docId },
  74. 'error processing update'
  75. )
  76. }
  77. }
  78. return cb()
  79. }
  80. )
  81. return RateLimiter.run(backgroundTask, callback)
  82. })
  83. },
  84. run() {
  85. if (Settings.shuttingDown) {
  86. return
  87. }
  88. return worker._waitForUpdateThenDispatchWorker(error => {
  89. if (error != null) {
  90. logger.error({ err: error }, 'Error in worker process')
  91. throw error
  92. } else {
  93. return worker.run()
  94. }
  95. })
  96. },
  97. }
  98. return worker
  99. },
  100. createAndStartDispatchers(number) {
  101. const RateLimiter = new RateLimitManager(number)
  102. _.times(number, function (shardNumber) {
  103. return DispatchManager.createDispatcher(RateLimiter, shardNumber).run()
  104. })
  105. },
  106. }