3_pack_docHistory_collection.coffee 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342
  1. Settings = require "settings-sharelatex"
  2. fs = require("fs")
  3. mongojs = require("mongojs")
  4. ObjectId = mongojs.ObjectId
  5. db = mongojs(Settings.mongo.url, ['docs','docHistory', 'docHistoryStats'])
  6. _ = require("underscore")
  7. async = require("async")
  8. exec = require("child_process").exec
  9. bson = require('bson')
  10. BSON = new bson()
  11. logger = {
  12. log: ->
  13. err: ->
  14. }
  15. needToExit = false
  16. handleExit = () ->
  17. needToExit = true
  18. console.log('Got signal. Shutting down.')
  19. process.on 'SIGINT', handleExit
  20. process.on 'SIGHUP', handleExit
  21. finished_docs_path = "/tmp/finished-docs-3"
  22. all_docs_path = "/tmp/all-docs-3"
  23. unmigrated_docs_path = "/tmp/unmigrated-docs-3"
  24. finished_docs = {}
  25. if fs.existsSync(finished_docs_path)
  26. for id in fs.readFileSync(finished_docs_path,'utf-8').split("\n")
  27. finished_docs[id] = true
  28. getAndWriteDocids = (callback)->
  29. console.log "finding all doc id's - #{new Date().toString()}"
  30. db.docs.find {}, {_id:1}, (err, ids)->
  31. console.log "total found docs in mongo #{ids.length} - #{new Date().toString()}"
  32. ids = _.pluck ids, '_id'
  33. ids = _.filter ids, (id)-> id?
  34. fileData = ids.join("\n")
  35. fs.writeFileSync all_docs_path + ".tmp", fileData
  36. fs.renameSync all_docs_path + ".tmp", all_docs_path
  37. callback(err, ids)
  38. loadDocIds = (callback)->
  39. console.log "loading doc ids from #{all_docs_path}"
  40. data = fs.readFileSync all_docs_path, "utf-8"
  41. ids = data.split("\n")
  42. console.log "loaded #{ids.length} doc ids from #{all_docs_path}"
  43. callback null, ids
  44. getDocIds = (callback)->
  45. exists = fs.existsSync all_docs_path
  46. if exists
  47. loadDocIds callback
  48. else
  49. getAndWriteDocids callback
  50. markDocAsProcessed = (doc_id, callback)->
  51. finished_docs[doc_id] = true
  52. fs.appendFile finished_docs_path, "#{doc_id}\n", callback
  53. markDocAsUnmigrated = (doc_id, callback)->
  54. console.log "#{doc_id} unmigrated"
  55. markDocAsProcessed doc_id, (err)->
  56. fs.appendFile unmigrated_docs_path, "#{doc_id}\n", callback
  57. checkIfDocHasBeenProccessed = (doc_id, callback)->
  58. callback(null, finished_docs[doc_id])
  59. processNext = (doc_id, callback)->
  60. if !doc_id? or doc_id.length == 0
  61. return callback()
  62. if needToExit
  63. return callback(new Error("graceful shutdown"))
  64. checkIfDocHasBeenProccessed doc_id, (err, hasBeenProcessed)->
  65. if hasBeenProcessed
  66. console.log "#{doc_id} already processed, skipping"
  67. return callback()
  68. PackManager._packDocHistory doc_id, {}, (err) ->
  69. if err?
  70. console.log "error processing #{doc_id}"
  71. markDocAsUnmigrated doc_id, callback
  72. else
  73. markDocAsProcessed doc_id, callback
  74. updateIndexes = (callback) ->
  75. async.series [
  76. (cb) ->
  77. console.log "create index"
  78. db.docHistory.ensureIndex { project_id: 1, "meta.end_ts": 1, "meta.start_ts": -1 }, { background: true }, cb
  79. (cb) ->
  80. console.log "drop index"
  81. db.docHistory.dropIndex { project_id: 1, "meta.end_ts": 1 }, cb
  82. (cb) ->
  83. console.log "drop index"
  84. db.docHistory.dropIndex { project_id: 1, "pack.0.meta.end_ts": 1, "meta.end_ts": 1}, cb
  85. ], (err, results) ->
  86. console.log "all done"
  87. callback(err)
  88. exports.migrate = (client, done = ->)->
  89. getDocIds (err, ids)->
  90. totalDocCount = ids.length
  91. alreadyFinishedCount = Object.keys(finished_docs).length
  92. t0 = Date.now()
  93. printProgress = () ->
  94. count = Object.keys(finished_docs).length
  95. processedFraction = (count-alreadyFinishedCount)/totalDocCount
  96. remainingFraction = (totalDocCount-count)/totalDocCount
  97. t = Date.now()
  98. dt = (t-t0)*remainingFraction/processedFraction
  99. estFinishTime = new Date(t + dt)
  100. console.log "completed #{count}/#{totalDocCount} processed=#{processedFraction.toFixed(2)} remaining=#{remainingFraction.toFixed(2)} elapsed=#{(t-t0)/1000} est Finish=#{estFinishTime}"
  101. interval = setInterval printProgress, 3*1000
  102. nextId = null
  103. testFn = () ->
  104. return false if needToExit
  105. id = ids.shift()
  106. while id? and finished_docs[id] # skip finished
  107. id = ids.shift()
  108. nextId = id
  109. return nextId?
  110. executeFn = (cb) ->
  111. processNext nextId, cb
  112. async.whilst testFn, executeFn, (err)->
  113. if err?
  114. console.error err, "at end of jobs"
  115. else
  116. console.log "finished at #{new Date}"
  117. clearInterval interval
  118. done(err)
  119. exports.rollback = (client, done)->
  120. done()
  121. # process.nextTick () ->
  122. # exports.migrate () ->
  123. # console.log "done"
  124. DAYS = 24 * 3600 * 1000 # one day in milliseconds
  125. # copied from track-changes/app/coffee/PackManager.coffee
  126. PackManager =
  127. MAX_SIZE: 1024*1024 # make these configurable parameters
  128. MAX_COUNT: 512
  129. convertDocsToPacks: (docs, callback) ->
  130. packs = []
  131. top = null
  132. docs.forEach (d,i) ->
  133. # skip existing packs
  134. if d.pack?
  135. top = null
  136. return
  137. sz = BSON.calculateObjectSize(d)
  138. # decide if this doc can be added to the current pack
  139. validLength = top? && (top.pack.length < PackManager.MAX_COUNT)
  140. validSize = top? && (top.sz + sz < PackManager.MAX_SIZE)
  141. bothPermanent = top? && (top.expiresAt? is false) && (d.expiresAt? is false)
  142. bothTemporary = top? && (top.expiresAt? is true) && (d.expiresAt? is true)
  143. within1Day = bothTemporary && (d.meta.start_ts - top.meta.start_ts < 24 * 3600 * 1000)
  144. if top? && validLength && validSize && (bothPermanent || (bothTemporary && within1Day))
  145. top.pack = top.pack.concat {v: d.v, meta: d.meta, op: d.op, _id: d._id}
  146. top.sz += sz
  147. top.n += 1
  148. top.v_end = d.v
  149. top.meta.end_ts = d.meta.end_ts
  150. top.expiresAt = d.expiresAt if top.expiresAt?
  151. return
  152. else
  153. # create a new pack
  154. top = _.clone(d)
  155. top.pack = [ {v: d.v, meta: d.meta, op: d.op, _id: d._id} ]
  156. top.meta = { start_ts: d.meta.start_ts, end_ts: d.meta.end_ts }
  157. top.sz = sz
  158. top.n = 1
  159. top.v_end = d.v
  160. delete top.op
  161. delete top._id
  162. packs.push top
  163. callback(null, packs)
  164. checkHistory: (docs, callback) ->
  165. errors = []
  166. prev = null
  167. error = (args...) ->
  168. errors.push args
  169. docs.forEach (d,i) ->
  170. if d.pack?
  171. n = d.pack.length
  172. last = d.pack[n-1]
  173. error('bad pack v_end', d) if d.v_end != last.v
  174. error('bad pack start_ts', d) if d.meta.start_ts != d.pack[0].meta.start_ts
  175. error('bad pack end_ts', d) if d.meta.end_ts != last.meta.end_ts
  176. d.pack.forEach (p, i) ->
  177. prev = v
  178. v = p.v
  179. error('bad version', v, 'in', p) if v <= prev
  180. #error('expired op', p, 'in pack') if p.expiresAt?
  181. else
  182. prev = v
  183. v = d.v
  184. error('bad version', v, 'in', d) if v <= prev
  185. if errors.length
  186. callback(errors)
  187. else
  188. callback()
  189. insertPack: (packObj, callback) ->
  190. bulk = db.docHistory.initializeOrderedBulkOp()
  191. doc_id = packObj.doc_id
  192. expect_nInserted = 1
  193. expect_nRemoved = packObj.pack.length
  194. logger.log {doc_id: doc_id}, "adding pack, removing #{expect_nRemoved} ops"
  195. bulk.insert packObj
  196. ids = (op._id for op in packObj.pack)
  197. bulk.find({_id:{$in:ids}}).remove()
  198. bulk.execute (err, result) ->
  199. if err?
  200. logger.error {doc_id: doc_id}, "error adding pack"
  201. callback(err, result)
  202. else if result.nInserted != expect_nInserted or result.nRemoved != expect_nRemoved
  203. logger.error {doc_id: doc_id, result}, "unexpected result adding pack"
  204. callback(new Error(
  205. msg: 'unexpected result'
  206. expected: {expect_nInserted, expect_nRemoved}
  207. ), result)
  208. else
  209. db.docHistoryStats.update {doc_id:doc_id}, {
  210. $inc:{update_count:-expect_nRemoved},
  211. $currentDate:{last_packed:true}
  212. }, {upsert:true}, () ->
  213. callback(err, result)
  214. # retrieve document ops/packs and check them
  215. getDocHistory: (doc_id, callback) ->
  216. db.docHistory.find({doc_id:ObjectId(doc_id)}).sort {v:1}, (err, docs) ->
  217. return callback(err) if err?
  218. # for safety, do a consistency check of the history
  219. logger.log {doc_id}, "checking history for document"
  220. PackManager.checkHistory docs, (err) ->
  221. return callback(err) if err?
  222. callback(err, docs)
  223. #PackManager.deleteExpiredPackOps docs, (err) ->
  224. # return callback(err) if err?
  225. # callback err, docs
  226. packDocHistory: (doc_id, options, callback) ->
  227. if typeof callback == "undefined" and typeof options == 'function'
  228. callback = options
  229. options = {}
  230. LockManager.runWithLock(
  231. "HistoryLock:#{doc_id}",
  232. (releaseLock) ->
  233. PackManager._packDocHistory(doc_id, options, releaseLock)
  234. , callback
  235. )
  236. _packDocHistory: (doc_id, options, callback) ->
  237. logger.log {doc_id},"starting pack operation for document history"
  238. PackManager.getDocHistory doc_id, (err, docs) ->
  239. return callback(err) if err?
  240. origDocs = 0
  241. origPacks = 0
  242. for d in docs
  243. if d.pack? then origPacks++ else origDocs++
  244. PackManager.convertDocsToPacks docs, (err, packs) ->
  245. return callback(err) if err?
  246. total = 0
  247. for p in packs
  248. total = total + p.pack.length
  249. logger.log {doc_id, origDocs, origPacks, newPacks: packs.length, totalOps: total}, "document stats"
  250. if packs.length
  251. if options['dry-run']
  252. logger.log {doc_id}, 'dry-run, skipping write packs'
  253. return callback()
  254. PackManager.savePacks packs, (err) ->
  255. return callback(err) if err?
  256. # check the history again
  257. PackManager.getDocHistory doc_id, callback
  258. else
  259. logger.log {doc_id}, "no packs to write"
  260. # keep a record that we checked this one to avoid rechecking it
  261. db.docHistoryStats.update {doc_id:doc_id}, {
  262. $currentDate:{last_checked:true}
  263. }, {upsert:true}, () ->
  264. callback null, null
  265. DB_WRITE_DELAY: 100
  266. savePacks: (packs, callback) ->
  267. async.eachSeries packs, PackManager.safeInsert, (err, result) ->
  268. if err?
  269. logger.log {err, result}, "error writing packs"
  270. callback err, result
  271. else
  272. callback()
  273. safeInsert: (packObj, callback) ->
  274. PackManager.insertPack packObj, (err, result) ->
  275. setTimeout () ->
  276. callback(err,result)
  277. , PackManager.DB_WRITE_DELAY
  278. deleteExpiredPackOps: (docs, callback) ->
  279. now = Date.now()
  280. toRemove = []
  281. toUpdate = []
  282. docs.forEach (d,i) ->
  283. if d.pack?
  284. newPack = d.pack.filter (op) ->
  285. if op.expiresAt? then op.expiresAt > now else true
  286. if newPack.length == 0
  287. toRemove.push d
  288. else if newPack.length < d.pack.length
  289. # adjust the pack properties
  290. d.pack = newPack
  291. first = d.pack[0]
  292. last = d.pack[d.pack.length - 1]
  293. d.v_end = last.v
  294. d.meta.start_ts = first.meta.start_ts
  295. d.meta.end_ts = last.meta.end_ts
  296. toUpdate.push d
  297. if toRemove.length or toUpdate.length
  298. bulk = db.docHistory.initializeOrderedBulkOp()
  299. toRemove.forEach (pack) ->
  300. console.log "would remove", pack
  301. #bulk.find({_id:pack._id}).removeOne()
  302. toUpdate.forEach (pack) ->
  303. console.log "would update", pack
  304. #bulk.find({_id:pack._id}).updateOne(pack);
  305. bulk.execute callback
  306. else
  307. callback()