app.js 8.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290
  1. const Metrics = require('@overleaf/metrics')
  2. const Settings = require('@overleaf/settings')
  3. Metrics.initialize(process.env.METRICS_APP_NAME || 'real-time')
  4. const async = require('async')
  5. const logger = require('logger-sharelatex')
  6. logger.initialize('real-time')
  7. Metrics.event_loop.monitor(logger)
  8. const express = require('express')
  9. const session = require('express-session')
  10. const redis = require('@overleaf/redis-wrapper')
  11. if (Settings.sentry && Settings.sentry.dsn) {
  12. logger.initializeErrorReporting(Settings.sentry.dsn)
  13. }
  14. const sessionRedisClient = redis.createClient(Settings.redis.websessions)
  15. const RedisStore = require('connect-redis')(session)
  16. const SessionSockets = require('./app/js/SessionSockets')
  17. const CookieParser = require('cookie-parser')
  18. const DrainManager = require('./app/js/DrainManager')
  19. const HealthCheckManager = require('./app/js/HealthCheckManager')
  20. const DeploymentManager = require('./app/js/DeploymentManager')
  21. // NOTE: debug is invoked for every blob that is put on the wire
  22. const socketIoLogger = {
  23. error(...message) {
  24. logger.debug({ fromSocketIo: true, originalLevel: 'error' }, ...message)
  25. },
  26. warn(...message) {
  27. logger.debug({ fromSocketIo: true, originalLevel: 'warn' }, ...message)
  28. },
  29. info() {},
  30. debug() {},
  31. log() {},
  32. }
  33. // monitor status file to take dark deployments out of the load-balancer
  34. DeploymentManager.initialise()
  35. // Set up socket.io server
  36. const app = express()
  37. const server = require('http').createServer(app)
  38. const io = require('socket.io').listen(server, {
  39. logger: socketIoLogger,
  40. })
  41. // Bind to sessions
  42. const sessionStore = new RedisStore({ client: sessionRedisClient })
  43. const cookieParser = CookieParser(Settings.security.sessionSecret)
  44. const sessionSockets = new SessionSockets(
  45. io,
  46. sessionStore,
  47. cookieParser,
  48. Settings.cookieName
  49. )
  50. Metrics.injectMetricsRoute(app)
  51. io.configure(function () {
  52. io.enable('browser client minification')
  53. io.enable('browser client etag')
  54. // Fix for Safari 5 error of "Error during WebSocket handshake: location mismatch"
  55. // See http://answers.dotcloud.com/question/578/problem-with-websocket-over-ssl-in-safari-with
  56. io.set('match origin protocol', true)
  57. // gzip uses a Node 0.8.x method of calling the gzip program which
  58. // doesn't work with 0.6.x
  59. // io.enable('browser client gzip')
  60. io.set('transports', [
  61. 'websocket',
  62. 'flashsocket',
  63. 'htmlfile',
  64. 'xhr-polling',
  65. 'jsonp-polling',
  66. ])
  67. })
  68. // a 200 response on '/' is required for load balancer health checks
  69. // these operate separately from kubernetes readiness checks
  70. app.get('/', function (req, res) {
  71. if (Settings.shutDownInProgress || DeploymentManager.deploymentIsClosed()) {
  72. res.sendStatus(503) // Service unavailable
  73. } else {
  74. res.send('real-time is open')
  75. }
  76. })
  77. app.get('/status', function (req, res) {
  78. if (Settings.shutDownInProgress) {
  79. res.sendStatus(503) // Service unavailable
  80. } else {
  81. res.send('real-time is alive')
  82. }
  83. })
  84. app.get('/debug/events', function (req, res) {
  85. Settings.debugEvents = parseInt(req.query.count, 10) || 20
  86. logger.info({ count: Settings.debugEvents }, 'starting debug mode')
  87. res.send(`debug mode will log next ${Settings.debugEvents} events`)
  88. })
  89. const rclient = require('@overleaf/redis-wrapper').createClient(
  90. Settings.redis.realtime
  91. )
  92. function healthCheck(req, res) {
  93. rclient.healthCheck(function (error) {
  94. if (error) {
  95. logger.err({ err: error }, 'failed redis health check')
  96. res.sendStatus(500)
  97. } else if (HealthCheckManager.isFailing()) {
  98. const status = HealthCheckManager.status()
  99. logger.err({ pubSubErrors: status }, 'failed pubsub health check')
  100. res.sendStatus(500)
  101. } else {
  102. res.sendStatus(200)
  103. }
  104. })
  105. }
  106. app.get(
  107. '/health_check',
  108. (req, res, next) => {
  109. if (Settings.shutDownComplete) {
  110. return res.sendStatus(503)
  111. }
  112. next()
  113. },
  114. healthCheck
  115. )
  116. app.get('/health_check/redis', healthCheck)
  117. // log http requests for routes defined from this point onwards
  118. app.use(Metrics.http.monitor(logger))
  119. const Router = require('./app/js/Router')
  120. Router.configure(app, io, sessionSockets)
  121. const WebsocketLoadBalancer = require('./app/js/WebsocketLoadBalancer')
  122. WebsocketLoadBalancer.listenForEditorEvents(io)
  123. const DocumentUpdaterController = require('./app/js/DocumentUpdaterController')
  124. DocumentUpdaterController.listenForUpdatesFromDocumentUpdater(io)
  125. const { port } = Settings.internal.realTime
  126. const { host } = Settings.internal.realTime
  127. server.listen(port, host, function (error) {
  128. if (error) {
  129. throw error
  130. }
  131. logger.info(`realtime starting up, listening on ${host}:${port}`)
  132. })
  133. // Stop huge stack traces in logs from all the socket.io parsing steps.
  134. Error.stackTraceLimit = 10
  135. function shutdownCleanly(signal) {
  136. const connectedClients = io.sockets.clients().length
  137. if (connectedClients === 0) {
  138. logger.info('no clients connected, exiting')
  139. process.exit()
  140. } else {
  141. logger.info(
  142. { connectedClients },
  143. 'clients still connected, not shutting down yet'
  144. )
  145. setTimeout(() => shutdownCleanly(signal), 30 * 1000)
  146. }
  147. }
  148. function drainAndShutdown(signal) {
  149. if (Settings.shutDownInProgress) {
  150. logger.info({ signal }, 'shutdown already in progress, ignoring signal')
  151. } else {
  152. Settings.shutDownInProgress = true
  153. const { statusCheckInterval } = Settings
  154. if (statusCheckInterval) {
  155. logger.info(
  156. { signal },
  157. `received interrupt, delay drain by ${statusCheckInterval}ms`
  158. )
  159. }
  160. setTimeout(function () {
  161. logger.info(
  162. { signal },
  163. `received interrupt, starting drain over ${shutdownDrainTimeWindow} mins`
  164. )
  165. DrainManager.startDrainTimeWindow(io, shutdownDrainTimeWindow, () => {
  166. setTimeout(() => {
  167. const staleClients = io.sockets.clients()
  168. if (staleClients.length !== 0) {
  169. logger.info(
  170. { staleClients: staleClients.map(client => client.id) },
  171. 'forcefully disconnecting stale clients'
  172. )
  173. staleClients.forEach(client => {
  174. client.disconnect()
  175. })
  176. }
  177. // Mark the node as unhealthy.
  178. Settings.shutDownComplete = true
  179. }, Settings.gracefulReconnectTimeoutMs)
  180. })
  181. shutdownCleanly(signal)
  182. }, statusCheckInterval)
  183. }
  184. }
  185. Settings.shutDownInProgress = false
  186. const shutdownDrainTimeWindow = parseInt(Settings.shutdownDrainTimeWindow, 10)
  187. if (Settings.shutdownDrainTimeWindow) {
  188. logger.info({ shutdownDrainTimeWindow }, 'shutdownDrainTimeWindow enabled')
  189. for (const signal of [
  190. 'SIGINT',
  191. 'SIGHUP',
  192. 'SIGQUIT',
  193. 'SIGUSR1',
  194. 'SIGUSR2',
  195. 'SIGTERM',
  196. 'SIGABRT',
  197. ]) {
  198. process.on(signal, drainAndShutdown)
  199. } // signal is passed as argument to event handler
  200. // global exception handler
  201. if (Settings.errors && Settings.errors.catchUncaughtErrors) {
  202. process.removeAllListeners('uncaughtException')
  203. process.on('uncaughtException', function (error) {
  204. if (
  205. [
  206. 'ETIMEDOUT',
  207. 'EHOSTUNREACH',
  208. 'EPIPE',
  209. 'ECONNRESET',
  210. 'ERR_STREAM_WRITE_AFTER_END',
  211. ].includes(error.code)
  212. ) {
  213. Metrics.inc('disconnected_write', 1, { status: error.code })
  214. return logger.warn(
  215. { err: error },
  216. 'attempted to write to disconnected client'
  217. )
  218. }
  219. logger.error({ err: error }, 'uncaught exception')
  220. if (Settings.errors && Settings.errors.shutdownOnUncaughtError) {
  221. drainAndShutdown('SIGABRT')
  222. }
  223. })
  224. }
  225. }
  226. if (Settings.continualPubsubTraffic) {
  227. logger.debug('continualPubsubTraffic enabled')
  228. const pubsubClient = redis.createClient(Settings.redis.pubsub)
  229. const clusterClient = redis.createClient(Settings.redis.websessions)
  230. const publishJob = function (channel, callback) {
  231. const checker = new HealthCheckManager(channel)
  232. logger.debug({ channel }, 'sending pub to keep connection alive')
  233. const json = JSON.stringify({
  234. health_check: true,
  235. key: checker.id,
  236. date: new Date().toString(),
  237. })
  238. Metrics.summary(`redis.publish.${channel}`, json.length)
  239. pubsubClient.publish(channel, json, function (err) {
  240. if (err) {
  241. logger.err({ err, channel }, 'error publishing pubsub traffic to redis')
  242. }
  243. const blob = JSON.stringify({ keep: 'alive' })
  244. Metrics.summary('redis.publish.cluster-continual-traffic', blob.length)
  245. clusterClient.publish('cluster-continual-traffic', blob, callback)
  246. })
  247. }
  248. const runPubSubTraffic = () =>
  249. async.map(['applied-ops', 'editor-events'], publishJob, () =>
  250. setTimeout(runPubSubTraffic, 1000 * 20)
  251. )
  252. runPubSubTraffic()
  253. }