PackManager.js 34 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251
  1. /* eslint-disable
  2. camelcase,
  3. no-unused-vars,
  4. */
  5. // TODO: This file was created by bulk-decaffeinate.
  6. // Fix any style issues and re-enable lint.
  7. /*
  8. * decaffeinate suggestions:
  9. * DS101: Remove unnecessary use of Array.from
  10. * DS102: Remove unnecessary code created because of implicit returns
  11. * DS205: Consider reworking code to avoid use of IIFEs
  12. * DS207: Consider shorter variations of null checks
  13. * Full docs: https://github.com/decaffeinate/decaffeinate/blob/master/docs/suggestions.md
  14. */
  15. let PackManager
  16. const async = require('async')
  17. const _ = require('underscore')
  18. const Bson = require('bson')
  19. const BSON = new Bson()
  20. const { db, ObjectId } = require('./mongodb')
  21. const logger = require('@overleaf/logger')
  22. const LockManager = require('./LockManager')
  23. const MongoAWS = require('./MongoAWS')
  24. const Metrics = require('@overleaf/metrics')
  25. const ProjectIterator = require('./ProjectIterator')
  26. const DocIterator = require('./DocIterator')
  27. const Settings = require('@overleaf/settings')
  28. const util = require('util')
  29. const keys = Settings.redis.lock.key_schema
  30. // Sharejs operations are stored in a 'pack' object
  31. //
  32. // e.g. a single sharejs update looks like
  33. //
  34. // {
  35. // "doc_id" : 549dae9e0a2a615c0c7f0c98,
  36. // "project_id" : 549dae9c0a2a615c0c7f0c8c,
  37. // "op" : [ {"p" : 6981, "d" : "?" } ],
  38. // "meta" : { "user_id" : 52933..., "start_ts" : 1422310693931, "end_ts" : 1422310693931 },
  39. // "v" : 17082
  40. // }
  41. //
  42. // and a pack looks like this
  43. //
  44. // {
  45. // "doc_id" : 549dae9e0a2a615c0c7f0c98,
  46. // "project_id" : 549dae9c0a2a615c0c7f0c8c,
  47. // "pack" : [ U1, U2, U3, ...., UN],
  48. // "meta" : { "user_id" : 52933..., "start_ts" : 1422310693931, "end_ts" : 1422310693931 },
  49. // "v" : 17082
  50. // "v_end" : ...
  51. // }
  52. //
  53. // where U1, U2, U3, .... are single updates stripped of their
  54. // doc_id and project_id fields (which are the same for all the
  55. // updates in the pack).
  56. //
  57. // The pack itself has v and meta fields, this makes it possible to
  58. // treat packs and single updates in a similar way.
  59. //
  60. // The v field of the pack itself is from the first entry U1, the
  61. // v_end field from UN. The meta.end_ts field of the pack itself is
  62. // from the last entry UN, the meta.start_ts field from U1.
  63. const DAYS = 24 * 3600 * 1000 // one day in milliseconds
  64. module.exports = PackManager = {
  65. MAX_SIZE: 1024 * 1024, // make these configurable parameters
  66. MAX_COUNT: 1024,
  67. insertCompressedUpdates(
  68. project_id,
  69. doc_id,
  70. lastUpdate,
  71. newUpdates,
  72. temporary,
  73. callback
  74. ) {
  75. if (callback == null) {
  76. callback = function () {}
  77. }
  78. if (newUpdates.length === 0) {
  79. return callback()
  80. }
  81. // never append permanent ops to a pack that will expire
  82. if (
  83. (lastUpdate != null ? lastUpdate.expiresAt : undefined) != null &&
  84. !temporary
  85. ) {
  86. lastUpdate = null
  87. }
  88. const updatesToFlush = []
  89. const updatesRemaining = newUpdates.slice()
  90. let n = (lastUpdate != null ? lastUpdate.n : undefined) || 0
  91. let sz = (lastUpdate != null ? lastUpdate.sz : undefined) || 0
  92. while (
  93. updatesRemaining.length &&
  94. n < PackManager.MAX_COUNT &&
  95. sz < PackManager.MAX_SIZE
  96. ) {
  97. const nextUpdate = updatesRemaining[0]
  98. const nextUpdateSize = BSON.calculateObjectSize(nextUpdate)
  99. if (nextUpdateSize + sz > PackManager.MAX_SIZE && n > 0) {
  100. break
  101. }
  102. n++
  103. sz += nextUpdateSize
  104. updatesToFlush.push(updatesRemaining.shift())
  105. }
  106. return PackManager.flushCompressedUpdates(
  107. project_id,
  108. doc_id,
  109. lastUpdate,
  110. updatesToFlush,
  111. temporary,
  112. function (error) {
  113. if (error != null) {
  114. return callback(error)
  115. }
  116. return PackManager.insertCompressedUpdates(
  117. project_id,
  118. doc_id,
  119. null,
  120. updatesRemaining,
  121. temporary,
  122. callback
  123. )
  124. }
  125. )
  126. },
  127. flushCompressedUpdates(
  128. project_id,
  129. doc_id,
  130. lastUpdate,
  131. newUpdates,
  132. temporary,
  133. callback
  134. ) {
  135. if (callback == null) {
  136. callback = function () {}
  137. }
  138. if (newUpdates.length === 0) {
  139. return callback()
  140. }
  141. let canAppend = false
  142. // check if it is safe to append to an existing pack
  143. if (lastUpdate != null) {
  144. if (!temporary && lastUpdate.expiresAt == null) {
  145. // permanent pack appends to permanent pack
  146. canAppend = true
  147. }
  148. const age =
  149. Date.now() -
  150. (lastUpdate.meta != null ? lastUpdate.meta.start_ts : undefined)
  151. if (temporary && lastUpdate.expiresAt != null && age < 1 * DAYS) {
  152. // temporary pack appends to temporary pack if same day
  153. canAppend = true
  154. }
  155. }
  156. if (canAppend) {
  157. return PackManager.appendUpdatesToExistingPack(
  158. project_id,
  159. doc_id,
  160. lastUpdate,
  161. newUpdates,
  162. temporary,
  163. callback
  164. )
  165. } else {
  166. return PackManager.insertUpdatesIntoNewPack(
  167. project_id,
  168. doc_id,
  169. newUpdates,
  170. temporary,
  171. callback
  172. )
  173. }
  174. },
  175. insertUpdatesIntoNewPack(
  176. project_id,
  177. doc_id,
  178. newUpdates,
  179. temporary,
  180. callback
  181. ) {
  182. if (callback == null) {
  183. callback = function () {}
  184. }
  185. const first = newUpdates[0]
  186. const last = newUpdates[newUpdates.length - 1]
  187. const n = newUpdates.length
  188. const sz = BSON.calculateObjectSize(newUpdates)
  189. const newPack = {
  190. project_id: ObjectId(project_id.toString()),
  191. doc_id: ObjectId(doc_id.toString()),
  192. pack: newUpdates,
  193. n,
  194. sz,
  195. meta: {
  196. start_ts: first.meta.start_ts,
  197. end_ts: last.meta.end_ts,
  198. },
  199. v: first.v,
  200. v_end: last.v,
  201. temporary,
  202. }
  203. if (temporary) {
  204. newPack.expiresAt = new Date(Date.now() + 7 * DAYS)
  205. newPack.last_checked = new Date(Date.now() + 30 * DAYS) // never check temporary packs
  206. }
  207. logger.debug(
  208. { project_id, doc_id, newUpdates },
  209. 'inserting updates into new pack'
  210. )
  211. return db.docHistory.insertOne(newPack, function (err) {
  212. if (err != null) {
  213. return callback(err)
  214. }
  215. Metrics.inc(`insert-pack-${temporary ? 'temporary' : 'permanent'}`)
  216. if (temporary) {
  217. return callback()
  218. } else {
  219. return PackManager.updateIndex(project_id, doc_id, callback)
  220. }
  221. })
  222. },
  223. appendUpdatesToExistingPack(
  224. project_id,
  225. doc_id,
  226. lastUpdate,
  227. newUpdates,
  228. temporary,
  229. callback
  230. ) {
  231. if (callback == null) {
  232. callback = function () {}
  233. }
  234. const first = newUpdates[0]
  235. const last = newUpdates[newUpdates.length - 1]
  236. const n = newUpdates.length
  237. const sz = BSON.calculateObjectSize(newUpdates)
  238. const query = {
  239. _id: lastUpdate._id,
  240. project_id: ObjectId(project_id.toString()),
  241. doc_id: ObjectId(doc_id.toString()),
  242. pack: { $exists: true },
  243. }
  244. const update = {
  245. $push: {
  246. pack: { $each: newUpdates },
  247. },
  248. $inc: {
  249. n: n,
  250. sz: sz,
  251. },
  252. $set: {
  253. 'meta.end_ts': last.meta.end_ts,
  254. v_end: last.v,
  255. },
  256. }
  257. if (lastUpdate.expiresAt && temporary) {
  258. update.$set.expiresAt = new Date(Date.now() + 7 * DAYS)
  259. }
  260. logger.debug(
  261. { project_id, doc_id, lastUpdate, newUpdates },
  262. 'appending updates to existing pack'
  263. )
  264. Metrics.inc(`append-pack-${temporary ? 'temporary' : 'permanent'}`)
  265. return db.docHistory.updateOne(query, update, callback)
  266. },
  267. // Retrieve all changes for a document
  268. getOpsByVersionRange(project_id, doc_id, fromVersion, toVersion, callback) {
  269. if (callback == null) {
  270. callback = function () {}
  271. }
  272. return PackManager.loadPacksByVersionRange(
  273. project_id,
  274. doc_id,
  275. fromVersion,
  276. toVersion,
  277. function (error) {
  278. if (error) return callback(error)
  279. const query = { doc_id: ObjectId(doc_id.toString()) }
  280. if (toVersion != null) {
  281. query.v = { $lte: toVersion }
  282. }
  283. if (fromVersion != null) {
  284. query.v_end = { $gte: fromVersion }
  285. }
  286. // console.log "query:", query
  287. return db.docHistory
  288. .find(query)
  289. .sort({ v: -1 })
  290. .toArray(function (err, result) {
  291. if (err != null) {
  292. return callback(err)
  293. }
  294. // console.log "getOpsByVersionRange:", err, result
  295. const updates = []
  296. const opInRange = function (op, from, to) {
  297. if (fromVersion != null && op.v < fromVersion) {
  298. return false
  299. }
  300. if (toVersion != null && op.v > toVersion) {
  301. return false
  302. }
  303. return true
  304. }
  305. for (const docHistory of Array.from(result)) {
  306. // console.log 'adding', docHistory.pack
  307. for (const op of Array.from(docHistory.pack.reverse())) {
  308. if (opInRange(op, fromVersion, toVersion)) {
  309. op.project_id = docHistory.project_id
  310. op.doc_id = docHistory.doc_id
  311. // console.log "added op", op.v, fromVersion, toVersion
  312. updates.push(op)
  313. }
  314. }
  315. }
  316. return callback(null, updates)
  317. })
  318. }
  319. )
  320. },
  321. loadPacksByVersionRange(
  322. project_id,
  323. doc_id,
  324. fromVersion,
  325. toVersion,
  326. callback
  327. ) {
  328. return PackManager.getIndex(doc_id, function (err, indexResult) {
  329. let pack
  330. if (err != null) {
  331. return callback(err)
  332. }
  333. const indexPacks =
  334. (indexResult != null ? indexResult.packs : undefined) || []
  335. const packInRange = function (pack, from, to) {
  336. if (fromVersion != null && pack.v_end < fromVersion) {
  337. return false
  338. }
  339. if (toVersion != null && pack.v > toVersion) {
  340. return false
  341. }
  342. return true
  343. }
  344. const neededIds = (() => {
  345. const result = []
  346. for (pack of Array.from(indexPacks)) {
  347. if (packInRange(pack, fromVersion, toVersion)) {
  348. result.push(pack._id)
  349. }
  350. }
  351. return result
  352. })()
  353. if (neededIds.length) {
  354. return PackManager.fetchPacksIfNeeded(
  355. project_id,
  356. doc_id,
  357. neededIds,
  358. callback
  359. )
  360. } else {
  361. return callback()
  362. }
  363. })
  364. },
  365. fetchPacksIfNeeded(project_id, doc_id, pack_ids, callback) {
  366. let id
  367. return db.docHistory
  368. .find(
  369. { _id: { $in: pack_ids.map(ObjectId) } },
  370. { projection: { _id: 1 } }
  371. )
  372. .toArray(function (err, loadedPacks) {
  373. if (err != null) {
  374. return callback(err)
  375. }
  376. const allPackIds = (() => {
  377. const result1 = []
  378. for (id of Array.from(pack_ids)) {
  379. result1.push(id.toString())
  380. }
  381. return result1
  382. })()
  383. const loadedPackIds = Array.from(loadedPacks).map(pack =>
  384. pack._id.toString()
  385. )
  386. const packIdsToFetch = _.difference(allPackIds, loadedPackIds)
  387. logger.debug(
  388. { project_id, doc_id, loadedPackIds, allPackIds, packIdsToFetch },
  389. 'analysed packs'
  390. )
  391. if (packIdsToFetch.length === 0) {
  392. return callback()
  393. }
  394. return async.eachLimit(
  395. packIdsToFetch,
  396. 4,
  397. (pack_id, cb) =>
  398. MongoAWS.unArchivePack(project_id, doc_id, pack_id, cb),
  399. function (err) {
  400. if (err != null) {
  401. return callback(err)
  402. }
  403. logger.debug({ project_id, doc_id }, 'done unarchiving')
  404. return callback()
  405. }
  406. )
  407. })
  408. },
  409. findAllDocsInProject(project_id, callback) {
  410. const docIdSet = new Set()
  411. async.series(
  412. [
  413. cb => {
  414. db.docHistory
  415. .find(
  416. { project_id: ObjectId(project_id) },
  417. { projection: { pack: false } }
  418. )
  419. .toArray((err, packs) => {
  420. if (err) return callback(err)
  421. packs.forEach(pack => {
  422. docIdSet.add(pack.doc_id.toString())
  423. })
  424. return cb()
  425. })
  426. },
  427. cb => {
  428. db.docHistoryIndex
  429. .find({ project_id: ObjectId(project_id) })
  430. .toArray((err, indexes) => {
  431. if (err) return callback(err)
  432. indexes.forEach(index => {
  433. docIdSet.add(index._id.toString())
  434. })
  435. return cb()
  436. })
  437. },
  438. ],
  439. err => {
  440. if (err) return callback(err)
  441. callback(null, [...docIdSet])
  442. }
  443. )
  444. },
  445. // rewrite any query using doc_id to use _id instead
  446. // (because docHistoryIndex uses the doc_id)
  447. _rewriteQueryForIndex(query) {
  448. const indexQuery = _.omit(query, 'doc_id')
  449. if ('doc_id' in query) {
  450. indexQuery._id = query.doc_id
  451. }
  452. return indexQuery
  453. },
  454. // Retrieve all changes across a project
  455. _findPacks(query, sortKeys, callback) {
  456. // get all the docHistory Entries
  457. return db.docHistory
  458. .find(query, { projection: { pack: false } })
  459. .sort(sortKeys)
  460. .toArray(function (err, packs) {
  461. let pack
  462. if (err != null) {
  463. return callback(err)
  464. }
  465. const allPacks = []
  466. const seenIds = {}
  467. for (pack of Array.from(packs)) {
  468. allPacks.push(pack)
  469. seenIds[pack._id] = true
  470. }
  471. const indexQuery = PackManager._rewriteQueryForIndex(query)
  472. return db.docHistoryIndex
  473. .find(indexQuery)
  474. .toArray(function (err, indexes) {
  475. if (err != null) {
  476. return callback(err)
  477. }
  478. for (const index of Array.from(indexes)) {
  479. for (pack of Array.from(index.packs)) {
  480. if (!seenIds[pack._id]) {
  481. pack.project_id = index.project_id
  482. pack.doc_id = index._id
  483. pack.fromIndex = true
  484. allPacks.push(pack)
  485. seenIds[pack._id] = true
  486. }
  487. }
  488. }
  489. return callback(null, allPacks)
  490. })
  491. })
  492. },
  493. makeProjectIterator(project_id, before, callback) {
  494. PackManager._findPacks(
  495. { project_id: ObjectId(project_id) },
  496. { 'meta.end_ts': -1 },
  497. function (err, allPacks) {
  498. if (err) return callback(err)
  499. callback(
  500. null,
  501. new ProjectIterator(allPacks, before, PackManager.getPackById)
  502. )
  503. }
  504. )
  505. },
  506. makeDocIterator(doc_id, callback) {
  507. PackManager._findPacks(
  508. { doc_id: ObjectId(doc_id) },
  509. { v: -1 },
  510. function (err, allPacks) {
  511. if (err) return callback(err)
  512. callback(null, new DocIterator(allPacks, PackManager.getPackById))
  513. }
  514. )
  515. },
  516. getPackById(project_id, doc_id, pack_id, callback) {
  517. return db.docHistory.findOne({ _id: pack_id }, function (err, pack) {
  518. if (err != null) {
  519. return callback(err)
  520. }
  521. if (pack == null) {
  522. return MongoAWS.unArchivePack(project_id, doc_id, pack_id, callback)
  523. } else if (pack.expiresAt != null && pack.temporary === false) {
  524. // we only need to touch the TTL when listing the changes in the project
  525. // because diffs on individual documents are always done after that
  526. return PackManager.increaseTTL(pack, callback)
  527. // only do this for cached packs, not temporary ones to avoid older packs
  528. // being kept longer than newer ones (which messes up the last update version)
  529. } else {
  530. return callback(null, pack)
  531. }
  532. })
  533. },
  534. increaseTTL(pack, callback) {
  535. if (pack.expiresAt < new Date(Date.now() + 6 * DAYS)) {
  536. // update cache expiry since we are using this pack
  537. return db.docHistory.updateOne(
  538. { _id: pack._id },
  539. { $set: { expiresAt: new Date(Date.now() + 7 * DAYS) } },
  540. err => callback(err, pack)
  541. )
  542. } else {
  543. return callback(null, pack)
  544. }
  545. },
  546. // Manage docHistoryIndex collection
  547. getIndex(doc_id, callback) {
  548. return db.docHistoryIndex.findOne(
  549. { _id: ObjectId(doc_id.toString()) },
  550. callback
  551. )
  552. },
  553. getPackFromIndex(doc_id, pack_id, callback) {
  554. return db.docHistoryIndex.findOne(
  555. { _id: ObjectId(doc_id.toString()), 'packs._id': pack_id },
  556. { projection: { 'packs.$': 1 } },
  557. callback
  558. )
  559. },
  560. getLastPackFromIndex(doc_id, callback) {
  561. return db.docHistoryIndex.findOne(
  562. { _id: ObjectId(doc_id.toString()) },
  563. { projection: { packs: { $slice: -1 } } },
  564. function (err, indexPack) {
  565. if (err != null) {
  566. return callback(err)
  567. }
  568. if (indexPack == null) {
  569. return callback()
  570. }
  571. return callback(null, indexPack[0])
  572. }
  573. )
  574. },
  575. getIndexWithKeys(doc_id, callback) {
  576. return PackManager.getIndex(doc_id, function (err, index) {
  577. if (err != null) {
  578. return callback(err)
  579. }
  580. if (index == null) {
  581. return callback()
  582. }
  583. for (const pack of Array.from(
  584. (index != null ? index.packs : undefined) || []
  585. )) {
  586. index[pack._id] = pack
  587. }
  588. return callback(null, index)
  589. })
  590. },
  591. initialiseIndex(project_id, doc_id, callback) {
  592. return PackManager.findCompletedPacks(
  593. project_id,
  594. doc_id,
  595. function (err, packs) {
  596. // console.log 'err', err, 'packs', packs, packs?.length
  597. if (err != null) {
  598. return callback(err)
  599. }
  600. if (packs == null) {
  601. return callback()
  602. }
  603. return PackManager.insertPacksIntoIndexWithLock(
  604. project_id,
  605. doc_id,
  606. packs,
  607. callback
  608. )
  609. }
  610. )
  611. },
  612. updateIndex(project_id, doc_id, callback) {
  613. // find all packs prior to current pack
  614. return PackManager.findUnindexedPacks(
  615. project_id,
  616. doc_id,
  617. function (err, newPacks) {
  618. if (err != null) {
  619. return callback(err)
  620. }
  621. if (newPacks == null || newPacks.length === 0) {
  622. return callback()
  623. }
  624. return PackManager.insertPacksIntoIndexWithLock(
  625. project_id,
  626. doc_id,
  627. newPacks,
  628. function (err) {
  629. if (err != null) {
  630. return callback(err)
  631. }
  632. logger.debug(
  633. { project_id, doc_id, newPacks },
  634. 'added new packs to index'
  635. )
  636. return callback()
  637. }
  638. )
  639. }
  640. )
  641. },
  642. findCompletedPacks(project_id, doc_id, callback) {
  643. const query = {
  644. doc_id: ObjectId(doc_id.toString()),
  645. expiresAt: { $exists: false },
  646. }
  647. return db.docHistory
  648. .find(query, { projection: { pack: false } })
  649. .sort({ v: 1 })
  650. .toArray(function (err, packs) {
  651. if (err != null) {
  652. return callback(err)
  653. }
  654. if (packs == null) {
  655. return callback()
  656. }
  657. if (!(packs != null ? packs.length : undefined)) {
  658. return callback()
  659. }
  660. const last = packs.pop() // discard the last pack, if it's still in progress
  661. if (last.finalised) {
  662. packs.push(last)
  663. } // it's finalised so we push it back to archive it
  664. return callback(null, packs)
  665. })
  666. },
  667. findPacks(project_id, doc_id, callback) {
  668. const query = {
  669. doc_id: ObjectId(doc_id.toString()),
  670. expiresAt: { $exists: false },
  671. }
  672. return db.docHistory
  673. .find(query, { projection: { pack: false } })
  674. .sort({ v: 1 })
  675. .toArray(function (err, packs) {
  676. if (err != null) {
  677. return callback(err)
  678. }
  679. if (packs == null) {
  680. return callback()
  681. }
  682. if (!(packs != null ? packs.length : undefined)) {
  683. return callback()
  684. }
  685. return callback(null, packs)
  686. })
  687. },
  688. findUnindexedPacks(project_id, doc_id, callback) {
  689. return PackManager.getIndexWithKeys(doc_id, function (err, indexResult) {
  690. if (err != null) {
  691. return callback(err)
  692. }
  693. return PackManager.findCompletedPacks(
  694. project_id,
  695. doc_id,
  696. function (err, historyPacks) {
  697. let pack
  698. if (err != null) {
  699. return callback(err)
  700. }
  701. if (historyPacks == null) {
  702. return callback()
  703. }
  704. // select only the new packs not already in the index
  705. let newPacks = (() => {
  706. const result = []
  707. for (pack of Array.from(historyPacks)) {
  708. if (
  709. (indexResult != null ? indexResult[pack._id] : undefined) ==
  710. null
  711. ) {
  712. result.push(pack)
  713. }
  714. }
  715. return result
  716. })()
  717. newPacks = (() => {
  718. const result1 = []
  719. for (pack of Array.from(newPacks)) {
  720. result1.push(
  721. _.omit(
  722. pack,
  723. 'doc_id',
  724. 'project_id',
  725. 'n',
  726. 'sz',
  727. 'last_checked',
  728. 'finalised'
  729. )
  730. )
  731. }
  732. return result1
  733. })()
  734. if (newPacks.length) {
  735. logger.debug(
  736. { project_id, doc_id, n: newPacks.length },
  737. 'found new packs'
  738. )
  739. }
  740. return callback(null, newPacks)
  741. }
  742. )
  743. })
  744. },
  745. insertPacksIntoIndexWithLock(project_id, doc_id, newPacks, callback) {
  746. return LockManager.runWithLock(
  747. keys.historyIndexLock({ doc_id }),
  748. releaseLock =>
  749. PackManager._insertPacksIntoIndex(
  750. project_id,
  751. doc_id,
  752. newPacks,
  753. releaseLock
  754. ),
  755. callback
  756. )
  757. },
  758. _insertPacksIntoIndex(project_id, doc_id, newPacks, callback) {
  759. return db.docHistoryIndex.updateOne(
  760. { _id: ObjectId(doc_id.toString()) },
  761. {
  762. $setOnInsert: { project_id: ObjectId(project_id.toString()) },
  763. $push: {
  764. packs: { $each: newPacks, $sort: { v: 1 } },
  765. },
  766. },
  767. {
  768. upsert: true,
  769. },
  770. callback
  771. )
  772. },
  773. // Archiving packs to S3
  774. archivePack(project_id, doc_id, pack_id, callback) {
  775. const clearFlagOnError = function (err, cb) {
  776. if (err != null) {
  777. // clear the inS3 flag on error
  778. return PackManager.clearPackAsArchiveInProgress(
  779. project_id,
  780. doc_id,
  781. pack_id,
  782. function (err2) {
  783. if (err2 != null) {
  784. return cb(err2)
  785. }
  786. return cb(err)
  787. }
  788. )
  789. } else {
  790. return cb()
  791. }
  792. }
  793. return async.series(
  794. [
  795. cb =>
  796. PackManager.checkArchiveNotInProgress(
  797. project_id,
  798. doc_id,
  799. pack_id,
  800. cb
  801. ),
  802. cb =>
  803. PackManager.markPackAsArchiveInProgress(
  804. project_id,
  805. doc_id,
  806. pack_id,
  807. cb
  808. ),
  809. cb =>
  810. MongoAWS.archivePack(project_id, doc_id, pack_id, err =>
  811. clearFlagOnError(err, cb)
  812. ),
  813. cb =>
  814. PackManager.checkArchivedPack(project_id, doc_id, pack_id, err =>
  815. clearFlagOnError(err, cb)
  816. ),
  817. cb => PackManager.markPackAsArchived(project_id, doc_id, pack_id, cb),
  818. cb =>
  819. PackManager.setTTLOnArchivedPack(
  820. project_id,
  821. doc_id,
  822. pack_id,
  823. callback
  824. ),
  825. ],
  826. callback
  827. )
  828. },
  829. checkArchivedPack(project_id, doc_id, pack_id, callback) {
  830. return db.docHistory.findOne({ _id: pack_id }, function (err, pack) {
  831. if (err != null) {
  832. return callback(err)
  833. }
  834. if (pack == null) {
  835. return callback(new Error('pack not found'))
  836. }
  837. return MongoAWS.readArchivedPack(
  838. project_id,
  839. doc_id,
  840. pack_id,
  841. function (err, result) {
  842. if (err) return callback(err)
  843. delete result.last_checked
  844. delete pack.last_checked
  845. // need to compare ids as ObjectIds with .equals()
  846. for (const key of ['_id', 'project_id', 'doc_id']) {
  847. if (result[key].equals(pack[key])) {
  848. result[key] = pack[key]
  849. }
  850. }
  851. for (let i = 0; i < result.pack.length; i++) {
  852. const op = result.pack[i]
  853. if (op._id != null && op._id.equals(pack.pack[i]._id)) {
  854. op._id = pack.pack[i]._id
  855. }
  856. }
  857. if (_.isEqual(pack, result)) {
  858. return callback()
  859. } else {
  860. logger.err(
  861. {
  862. pack,
  863. result,
  864. jsondiff: JSON.stringify(pack) === JSON.stringify(result),
  865. },
  866. 'difference when comparing packs'
  867. )
  868. return callback(
  869. new Error('pack retrieved from s3 does not match pack in mongo')
  870. )
  871. }
  872. }
  873. )
  874. })
  875. },
  876. // Extra methods to test archive/unarchive for a doc_id
  877. pushOldPacks(project_id, doc_id, callback) {
  878. return PackManager.findPacks(project_id, doc_id, function (err, packs) {
  879. if (err != null) {
  880. return callback(err)
  881. }
  882. if (!(packs != null ? packs.length : undefined)) {
  883. return callback()
  884. }
  885. return PackManager.processOldPack(
  886. project_id,
  887. doc_id,
  888. packs[0]._id,
  889. callback
  890. )
  891. })
  892. },
  893. pullOldPacks(project_id, doc_id, callback) {
  894. return PackManager.loadPacksByVersionRange(
  895. project_id,
  896. doc_id,
  897. null,
  898. null,
  899. callback
  900. )
  901. },
  902. // Processing old packs via worker
  903. processOldPack(project_id, doc_id, pack_id, callback) {
  904. const markAsChecked = err =>
  905. PackManager.markPackAsChecked(
  906. project_id,
  907. doc_id,
  908. pack_id,
  909. function (err2) {
  910. if (err2 != null) {
  911. return callback(err2)
  912. }
  913. return callback(err)
  914. }
  915. )
  916. logger.debug({ project_id, doc_id }, 'processing old packs')
  917. return db.docHistory.findOne({ _id: pack_id }, function (err, pack) {
  918. if (err != null) {
  919. return markAsChecked(err)
  920. }
  921. if (pack == null) {
  922. return markAsChecked()
  923. }
  924. if (pack.expiresAt != null) {
  925. return callback()
  926. } // return directly
  927. return PackManager.finaliseIfNeeded(
  928. project_id,
  929. doc_id,
  930. pack._id,
  931. pack,
  932. function (err) {
  933. if (err != null) {
  934. return markAsChecked(err)
  935. }
  936. return PackManager.updateIndexIfNeeded(
  937. project_id,
  938. doc_id,
  939. function (err) {
  940. if (err != null) {
  941. return markAsChecked(err)
  942. }
  943. return PackManager.findUnarchivedPacks(
  944. project_id,
  945. doc_id,
  946. function (err, unarchivedPacks) {
  947. if (err != null) {
  948. return markAsChecked(err)
  949. }
  950. if (
  951. !(unarchivedPacks != null
  952. ? unarchivedPacks.length
  953. : undefined)
  954. ) {
  955. logger.debug(
  956. { project_id, doc_id },
  957. 'no packs need archiving'
  958. )
  959. return markAsChecked()
  960. }
  961. return async.eachSeries(
  962. unarchivedPacks,
  963. (pack, cb) =>
  964. PackManager.archivePack(project_id, doc_id, pack._id, cb),
  965. function (err) {
  966. if (err != null) {
  967. return markAsChecked(err)
  968. }
  969. logger.debug({ project_id, doc_id }, 'done processing')
  970. return markAsChecked()
  971. }
  972. )
  973. }
  974. )
  975. }
  976. )
  977. }
  978. )
  979. })
  980. },
  981. finaliseIfNeeded(project_id, doc_id, pack_id, pack, callback) {
  982. const sz = pack.sz / (1024 * 1024) // in fractions of a megabyte
  983. const n = pack.n / 1024 // in fraction of 1024 ops
  984. const age = (Date.now() - pack.meta.end_ts) / DAYS
  985. if (age < 30) {
  986. // always keep if less than 1 month old
  987. logger.debug(
  988. { project_id, doc_id, pack_id, age },
  989. 'less than 30 days old'
  990. )
  991. return callback()
  992. }
  993. // compute an archiving threshold which decreases for each month of age
  994. const archive_threshold = 30 / age
  995. if (sz > archive_threshold || n > archive_threshold || age > 90) {
  996. logger.debug(
  997. { project_id, doc_id, pack_id, age, archive_threshold, sz, n },
  998. 'meets archive threshold'
  999. )
  1000. return PackManager.markPackAsFinalisedWithLock(
  1001. project_id,
  1002. doc_id,
  1003. pack_id,
  1004. callback
  1005. )
  1006. } else {
  1007. logger.debug(
  1008. { project_id, doc_id, pack_id, age, archive_threshold, sz, n },
  1009. 'does not meet archive threshold'
  1010. )
  1011. return callback()
  1012. }
  1013. },
  1014. markPackAsFinalisedWithLock(project_id, doc_id, pack_id, callback) {
  1015. return LockManager.runWithLock(
  1016. keys.historyLock({ doc_id }),
  1017. releaseLock =>
  1018. PackManager._markPackAsFinalised(
  1019. project_id,
  1020. doc_id,
  1021. pack_id,
  1022. releaseLock
  1023. ),
  1024. callback
  1025. )
  1026. },
  1027. _markPackAsFinalised(project_id, doc_id, pack_id, callback) {
  1028. logger.debug({ project_id, doc_id, pack_id }, 'marking pack as finalised')
  1029. return db.docHistory.updateOne(
  1030. { _id: pack_id },
  1031. { $set: { finalised: true } },
  1032. callback
  1033. )
  1034. },
  1035. updateIndexIfNeeded(project_id, doc_id, callback) {
  1036. logger.debug({ project_id, doc_id }, 'archiving old packs')
  1037. return PackManager.getIndexWithKeys(doc_id, function (err, index) {
  1038. if (err != null) {
  1039. return callback(err)
  1040. }
  1041. if (index == null) {
  1042. return PackManager.initialiseIndex(project_id, doc_id, callback)
  1043. } else {
  1044. return PackManager.updateIndex(project_id, doc_id, callback)
  1045. }
  1046. })
  1047. },
  1048. markPackAsChecked(project_id, doc_id, pack_id, callback) {
  1049. logger.debug({ project_id, doc_id, pack_id }, 'marking pack as checked')
  1050. return db.docHistory.updateOne(
  1051. { _id: pack_id },
  1052. { $currentDate: { last_checked: true } },
  1053. callback
  1054. )
  1055. },
  1056. findUnarchivedPacks(project_id, doc_id, callback) {
  1057. return PackManager.getIndex(doc_id, function (err, indexResult) {
  1058. if (err != null) {
  1059. return callback(err)
  1060. }
  1061. const indexPacks =
  1062. (indexResult != null ? indexResult.packs : undefined) || []
  1063. const unArchivedPacks = (() => {
  1064. const result = []
  1065. for (const pack of Array.from(indexPacks)) {
  1066. if (pack.inS3 == null) {
  1067. result.push(pack)
  1068. }
  1069. }
  1070. return result
  1071. })()
  1072. if (unArchivedPacks.length) {
  1073. logger.debug(
  1074. { project_id, doc_id, n: unArchivedPacks.length },
  1075. 'find unarchived packs'
  1076. )
  1077. }
  1078. return callback(null, unArchivedPacks)
  1079. })
  1080. },
  1081. // Archive locking flags
  1082. checkArchiveNotInProgress(project_id, doc_id, pack_id, callback) {
  1083. logger.debug(
  1084. { project_id, doc_id, pack_id },
  1085. 'checking if archive in progress'
  1086. )
  1087. return PackManager.getPackFromIndex(
  1088. doc_id,
  1089. pack_id,
  1090. function (err, result) {
  1091. if (err != null) {
  1092. return callback(err)
  1093. }
  1094. if (result == null) {
  1095. return callback(new Error('pack not found in index'))
  1096. }
  1097. if (result.inS3) {
  1098. return callback(new Error('pack archiving already done'))
  1099. } else if (result.inS3 != null) {
  1100. return callback(new Error('pack archiving already in progress'))
  1101. } else {
  1102. return callback()
  1103. }
  1104. }
  1105. )
  1106. },
  1107. markPackAsArchiveInProgress(project_id, doc_id, pack_id, callback) {
  1108. logger.debug(
  1109. { project_id, doc_id },
  1110. 'marking pack as archive in progress status'
  1111. )
  1112. return db.docHistoryIndex.findOneAndUpdate(
  1113. {
  1114. _id: ObjectId(doc_id.toString()),
  1115. packs: { $elemMatch: { _id: pack_id, inS3: { $exists: false } } },
  1116. },
  1117. { $set: { 'packs.$.inS3': false } },
  1118. { projection: { 'packs.$': 1 } },
  1119. function (err, result) {
  1120. if (err != null) {
  1121. return callback(err)
  1122. }
  1123. if (!result.value) {
  1124. return callback(new Error('archive is already in progress'))
  1125. }
  1126. logger.debug(
  1127. { project_id, doc_id, pack_id },
  1128. 'marked as archive in progress'
  1129. )
  1130. return callback()
  1131. }
  1132. )
  1133. },
  1134. clearPackAsArchiveInProgress(project_id, doc_id, pack_id, callback) {
  1135. logger.debug(
  1136. { project_id, doc_id, pack_id },
  1137. 'clearing as archive in progress'
  1138. )
  1139. return db.docHistoryIndex.updateOne(
  1140. {
  1141. _id: ObjectId(doc_id.toString()),
  1142. packs: { $elemMatch: { _id: pack_id, inS3: false } },
  1143. },
  1144. { $unset: { 'packs.$.inS3': true } },
  1145. callback
  1146. )
  1147. },
  1148. markPackAsArchived(project_id, doc_id, pack_id, callback) {
  1149. logger.debug({ project_id, doc_id, pack_id }, 'marking pack as archived')
  1150. return db.docHistoryIndex.findOneAndUpdate(
  1151. {
  1152. _id: ObjectId(doc_id.toString()),
  1153. packs: { $elemMatch: { _id: pack_id, inS3: false } },
  1154. },
  1155. { $set: { 'packs.$.inS3': true } },
  1156. { projection: { 'packs.$': 1 } },
  1157. function (err, result) {
  1158. if (err != null) {
  1159. return callback(err)
  1160. }
  1161. if (!result.value) {
  1162. return callback(new Error('archive is not marked as progress'))
  1163. }
  1164. logger.debug({ project_id, doc_id, pack_id }, 'marked as archived')
  1165. return callback()
  1166. }
  1167. )
  1168. },
  1169. setTTLOnArchivedPack(project_id, doc_id, pack_id, callback) {
  1170. return db.docHistory.updateOne(
  1171. { _id: pack_id },
  1172. { $set: { expiresAt: new Date(Date.now() + 1 * DAYS) } },
  1173. function (err) {
  1174. if (err) {
  1175. return callback(err)
  1176. }
  1177. logger.debug({ project_id, doc_id, pack_id }, 'set expiry on pack')
  1178. return callback()
  1179. }
  1180. )
  1181. },
  1182. }
  1183. module.exports.promises = {
  1184. getOpsByVersionRange: util.promisify(PackManager.getOpsByVersionRange),
  1185. findAllDocsInProject: util.promisify(PackManager.findAllDocsInProject),
  1186. makeDocIterator: util.promisify(PackManager.makeDocIterator),
  1187. }
  1188. // _getOneDayInFutureWithRandomDelay: ->
  1189. // thirtyMins = 1000 * 60 * 30
  1190. // randomThirtyMinMax = Math.ceil(Math.random() * thirtyMins)
  1191. // return new Date(Date.now() + randomThirtyMinMax + 1*DAYS)