DocumentUpdaterController.js 5.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185
  1. import logger from '@overleaf/logger'
  2. import settings from '@overleaf/settings'
  3. import RedisClientManager from './RedisClientManager.js'
  4. import SafeJsonParse from './SafeJsonParse.js'
  5. import EventLogger from './EventLogger.js'
  6. import HealthCheckManager from './HealthCheckManager.js'
  7. import RoomManager from './RoomManager.js'
  8. import ChannelManager from './ChannelManager.js'
  9. import metrics from '@overleaf/metrics'
  10. let DocumentUpdaterController
  11. export default DocumentUpdaterController = {
  12. // DocumentUpdaterController is responsible for updates that come via Redis
  13. // Pub/Sub from the document updater.
  14. rclientList: RedisClientManager.createClientList(settings.redis.pubsub),
  15. listenForUpdatesFromDocumentUpdater(io) {
  16. logger.debug(
  17. { rclients: this.rclientList.length },
  18. 'listening for applied-ops events'
  19. )
  20. for (const rclient of this.rclientList) {
  21. rclient.subscribe('applied-ops')
  22. rclient.on('message', function (channel, message) {
  23. metrics.inc('rclient', 0.001) // global event rate metric
  24. if (settings.debugEvents > 0) {
  25. EventLogger.debugEvent(channel, message)
  26. }
  27. DocumentUpdaterController._processMessageFromDocumentUpdater(
  28. io,
  29. channel,
  30. message
  31. )
  32. })
  33. }
  34. // create metrics for each redis instance only when we have multiple redis clients
  35. if (this.rclientList.length > 1) {
  36. this.rclientList.forEach((rclient, i) => {
  37. // per client event rate metric
  38. const metricName = `rclient-${i}`
  39. rclient.on('message', () => metrics.inc(metricName, 0.001))
  40. })
  41. }
  42. this.handleRoomUpdates(this.rclientList)
  43. },
  44. handleRoomUpdates(rclientSubList) {
  45. const roomEvents = RoomManager.eventSource()
  46. roomEvents.on('doc-active', function (docId) {
  47. const subscribePromises = rclientSubList.map(rclient =>
  48. ChannelManager.subscribe(rclient, 'applied-ops', docId)
  49. )
  50. RoomManager.emitOnCompletion(subscribePromises, `doc-subscribed-${docId}`)
  51. })
  52. roomEvents.on('doc-empty', docId =>
  53. rclientSubList.map(rclient =>
  54. ChannelManager.unsubscribe(rclient, 'applied-ops', docId)
  55. )
  56. )
  57. },
  58. _processMessageFromDocumentUpdater(io, channel, message) {
  59. SafeJsonParse.parse(message, function (error, message) {
  60. if (error) {
  61. logger.error({ err: error, channel }, 'error parsing JSON')
  62. return
  63. }
  64. if (message.op) {
  65. if (message._id && settings.checkEventOrder) {
  66. const status = EventLogger.checkEventOrder(
  67. 'applied-ops',
  68. message._id,
  69. message
  70. )
  71. if (status === 'duplicate') {
  72. return // skip duplicate events
  73. }
  74. }
  75. DocumentUpdaterController._applyUpdateFromDocumentUpdater(
  76. io,
  77. message.doc_id,
  78. message.op
  79. )
  80. } else if (message.error) {
  81. DocumentUpdaterController._processErrorFromDocumentUpdater(
  82. io,
  83. message.doc_id,
  84. message.error,
  85. message
  86. )
  87. } else if (message.health_check) {
  88. logger.debug(
  89. { message },
  90. 'got health check message in applied ops channel'
  91. )
  92. HealthCheckManager.check(channel, message.key)
  93. }
  94. })
  95. },
  96. _applyUpdateFromDocumentUpdater(io, docId, update) {
  97. let client
  98. const clientList = io.sockets.clients(docId)
  99. // avoid unnecessary work if no clients are connected
  100. if (clientList.length === 0) {
  101. return
  102. }
  103. update.meta = update.meta || {}
  104. const { tsRT: realTimeIngestionTime } = update.meta
  105. delete update.meta.tsRT
  106. // send updates to clients
  107. logger.debug(
  108. {
  109. docId,
  110. version: update.v,
  111. source: update.meta && update.meta.source,
  112. socketIoClients: clientList.map(client => client.id),
  113. },
  114. 'distributing updates to clients'
  115. )
  116. const seen = {}
  117. // send messages only to unique clients (due to duplicate entries in io.sockets.clients)
  118. for (client of clientList) {
  119. if (!seen[client.id]) {
  120. seen[client.id] = true
  121. if (client.publicId === update.meta.source) {
  122. logger.debug(
  123. {
  124. docId,
  125. version: update.v,
  126. source: update.meta.source,
  127. },
  128. 'distributing update to sender'
  129. )
  130. metrics.histogram(
  131. 'update-processing-time',
  132. performance.now() - realTimeIngestionTime,
  133. [
  134. 0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 20, 50, 100, 200, 500, 1000,
  135. 2000, 5000, 10000,
  136. ],
  137. { path: 'sharejs' }
  138. )
  139. client.emit('otUpdateApplied', { v: update.v, doc: update.doc })
  140. } else if (!update.dup) {
  141. // Duplicate ops should just be sent back to sending client for acknowledgement
  142. logger.debug(
  143. {
  144. docId,
  145. version: update.v,
  146. source: update.meta.source,
  147. clientId: client.id,
  148. },
  149. 'distributing update to collaborator'
  150. )
  151. client.emit('otUpdateApplied', update)
  152. }
  153. }
  154. }
  155. if (Object.keys(seen).length < clientList.length) {
  156. metrics.inc('socket-io.duplicate-clients', 0.1)
  157. logger.debug(
  158. {
  159. docId,
  160. socketIoClients: clientList.map(client => client.id),
  161. },
  162. 'discarded duplicate clients'
  163. )
  164. }
  165. },
  166. _processErrorFromDocumentUpdater(io, docId, error, message) {
  167. for (const client of io.sockets.clients(docId)) {
  168. logger.warn(
  169. { err: error, docId, clientId: client.id },
  170. 'error from document updater, disconnecting client'
  171. )
  172. client.emit('otUpdateError', error, message)
  173. client.disconnect()
  174. }
  175. },
  176. }