TpdsWorker.coffee 3.1 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798
  1. async = require('async')
  2. request = require('request')
  3. keys = require('./app/js/infrastructure/Keys')
  4. settings = require('settings-sharelatex')
  5. logger = require('logger-sharelatex')
  6. _ = require('underscore')
  7. childProcess = require("child_process")
  8. metrics = require("./app/js/infrastructure/Metrics")
  9. fiveMinutes = 5 * 60 * 1000
  10. processingFuncs =
  11. sendDoc : (options, callback)->
  12. if !options.docLines? || options.docLines.length == 0
  13. logger.err options:options, "doc lines not added to options for processing"
  14. return callback()
  15. docLines = options.docLines.reduce (singleLine, line)-> "#{singleLine}\n#{line}"
  16. post = request(options)
  17. post.on 'error', (err)->
  18. if err?
  19. callback(err)
  20. else
  21. callback()
  22. post.on 'end', callback
  23. post.write(docLines, 'utf-8')
  24. standardHttpRequest: (options, callback)->
  25. request options, (err, reponse, body)->
  26. if err?
  27. callback(err)
  28. else
  29. callback()
  30. pipeStreamFrom: (options, callback)->
  31. if options.filePath == "/droppy/main.tex"
  32. request options.streamOrigin, (err,res, body)->
  33. logger.log options:options, body:body
  34. origin = request(options.streamOrigin)
  35. origin.on 'error', (err)->
  36. logger.error err:err, options:options, "something went wrong in pipeStreamFrom origin"
  37. if err?
  38. callback(err)
  39. else
  40. callback()
  41. dest = request(options)
  42. origin.pipe(dest)
  43. dest.on "error", (err)->
  44. logger.error err:err, options:options, "something went wrong in pipeStreamFrom dest"
  45. if err?
  46. callback(err)
  47. else
  48. callback()
  49. dest.on 'end', callback
  50. workerRegistration = (groupKey, method, options, callback)->
  51. callback = _.once callback
  52. setTimeout callback, fiveMinutes
  53. metrics.inc "tpds-worker-processing"
  54. logger.log groupKey:groupKey, method:method, options:options, "processing http request from queue"
  55. processingFuncs[method] options, (err)->
  56. if err?
  57. logger.err err:err, user_id:groupKey, method:method, options:options, "something went wrong processing tpdsUpdateSender update"
  58. return callback("skip-after-retry")
  59. callback()
  60. setupWebToTpdsWorkers = (queueName)->
  61. logger.log worker_count:worker_count, queueName:queueName, "fairy workers"
  62. worker_count = 4
  63. while worker_count-- > 0
  64. workerQueueRef = require('fairy').connect(settings.redis.fairy).queue(queueName)
  65. workerQueueRef.polling_interval = 100
  66. workerQueueRef.regist workerRegistration
  67. cleanupPreviousQueues = (queueName, callback)->
  68. #cleanup queues then setup workers
  69. fairy = require('fairy').connect(settings.redis.fairy)
  70. queuePrefix = "FAIRY:QUEUED:#{queueName}:"
  71. fairy.redis.keys "#{queuePrefix}*", (err, keys)->
  72. logger.log "#{keys.length} fairy queues need cleanup"
  73. queueNames = keys.map (key)->
  74. key.replace queuePrefix, ""
  75. cleanupJobs = queueNames.map (projectQueueName)->
  76. return (cb)->
  77. cleanup = childProcess.fork(__dirname + '/cleanup.js', [queueName, projectQueueName])
  78. cleanup.on 'exit', cb
  79. async.series cleanupJobs, callback
  80. cleanupPreviousQueues keys.queue.web_to_tpds_http_requests, ->
  81. setupWebToTpdsWorkers keys.queue.web_to_tpds_http_requests
  82. cleanupPreviousQueues keys.queue.tpds_to_web_http_requests, ->
  83. setupWebToTpdsWorkers keys.queue.tpds_to_web_http_requests