app.js 9.0 KB

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