GcsPersistor.js 9.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333
  1. const fs = require('fs')
  2. const { pipeline } = require('stream/promises')
  3. const { PassThrough } = require('stream')
  4. const { Storage, IdempotencyStrategy } = require('@google-cloud/storage')
  5. const { WriteError, ReadError, NotFoundError } = require('./Errors')
  6. const asyncPool = require('tiny-async-pool')
  7. const AbstractPersistor = require('./AbstractPersistor')
  8. const PersistorHelper = require('./PersistorHelper')
  9. const Logger = require('@overleaf/logger')
  10. module.exports = class GcsPersistor extends AbstractPersistor {
  11. constructor(settings) {
  12. super()
  13. this.settings = settings
  14. // endpoint settings will be null by default except for tests
  15. // that's OK - GCS uses the locally-configured service account by default
  16. const storageOptions = {}
  17. if (this.settings.endpoint) {
  18. storageOptions.projectId = this.settings.endpoint.projectId
  19. storageOptions.apiEndpoint = this.settings.endpoint.apiEndpoint
  20. }
  21. storageOptions.retryOptions = { ...this.settings.retryOptions }
  22. if (storageOptions.retryOptions) {
  23. if (storageOptions.retryOptions.idempotencyStrategy) {
  24. const value =
  25. IdempotencyStrategy[this.settings.retryOptions.idempotencyStrategy]
  26. if (value === undefined) {
  27. throw new Error(
  28. 'Unrecognised value for retryOptions.idempotencyStrategy'
  29. )
  30. }
  31. Logger.info(
  32. `Setting retryOptions.idempotencyStrategy to ${storageOptions.retryOptions.idempotencyStrategy} (${value})`
  33. )
  34. storageOptions.retryOptions.idempotencyStrategy = value
  35. }
  36. }
  37. this.storage = new Storage(storageOptions)
  38. }
  39. async sendFile(bucketName, key, fsPath) {
  40. return await this.sendStream(bucketName, key, fs.createReadStream(fsPath))
  41. }
  42. async sendStream(bucketName, key, readStream, opts = {}) {
  43. try {
  44. // egress from us to gcs
  45. const observeOptions = {
  46. metric: 'gcs.egress',
  47. Metrics: this.settings.Metrics,
  48. }
  49. let sourceMd5 = opts.sourceMd5
  50. if (!sourceMd5) {
  51. // if there is no supplied md5 hash, we calculate the hash as the data passes through
  52. observeOptions.hash = 'md5'
  53. }
  54. const observer = new PersistorHelper.ObserverStream(observeOptions)
  55. const writeOptions = {
  56. // disabling of resumable uploads is recommended by Google:
  57. resumable: false,
  58. }
  59. if (sourceMd5) {
  60. writeOptions.validation = 'md5'
  61. writeOptions.metadata = writeOptions.metadata || {}
  62. writeOptions.metadata.md5Hash = PersistorHelper.hexToBase64(sourceMd5)
  63. }
  64. if (opts.contentType) {
  65. writeOptions.metadata = writeOptions.metadata || {}
  66. writeOptions.metadata.contentType = opts.contentType
  67. }
  68. if (opts.contentEncoding) {
  69. writeOptions.metadata = writeOptions.metadata || {}
  70. writeOptions.metadata.contentEncoding = opts.contentEncoding
  71. }
  72. const uploadStream = this.storage
  73. .bucket(bucketName)
  74. .file(key)
  75. .createWriteStream(writeOptions)
  76. await pipeline(readStream, observer, uploadStream)
  77. // if we didn't have an md5 hash, we should compare our computed one with Google's
  78. // as we couldn't tell GCS about it beforehand
  79. if (!sourceMd5) {
  80. sourceMd5 = observer.getHash()
  81. // throws on mismatch
  82. await PersistorHelper.verifyMd5(this, bucketName, key, sourceMd5)
  83. }
  84. } catch (err) {
  85. throw PersistorHelper.wrapError(
  86. err,
  87. 'upload to GCS failed',
  88. { bucketName, key },
  89. WriteError
  90. )
  91. }
  92. }
  93. async getObjectStream(bucketName, key, opts = {}) {
  94. const stream = this.storage
  95. .bucket(bucketName)
  96. .file(key)
  97. .createReadStream({ decompress: false, ...opts })
  98. try {
  99. await new Promise((resolve, reject) => {
  100. stream.on('response', res => {
  101. switch (res.statusCode) {
  102. case 200: // full response
  103. case 206: // partial response
  104. return resolve()
  105. case 404:
  106. return reject(new NotFoundError())
  107. default:
  108. return reject(new Error('non success status: ' + res.statusCode))
  109. }
  110. })
  111. stream.on('error', reject)
  112. stream.read(0) // kick off request
  113. })
  114. } catch (err) {
  115. throw PersistorHelper.wrapError(
  116. err,
  117. 'error reading file from GCS',
  118. { bucketName, key, opts },
  119. ReadError
  120. )
  121. }
  122. // ingress to us from gcs
  123. const observer = new PersistorHelper.ObserverStream({
  124. metric: 'gcs.ingress',
  125. Metrics: this.settings.Metrics,
  126. })
  127. const pass = new PassThrough()
  128. pipeline(stream, observer, pass).catch(() => {})
  129. return pass
  130. }
  131. async getRedirectUrl(bucketName, key) {
  132. if (this.settings.unsignedUrls) {
  133. // Construct a direct URL to the object download endpoint
  134. // (see https://cloud.google.com/storage/docs/request-endpoints#json-api)
  135. const apiEndpoint =
  136. this.settings.endpoint.apiEndpoint || 'https://storage.googleapis.com'
  137. return `${apiEndpoint}/download/storage/v1/b/${bucketName}/o/${key}?alt=media`
  138. }
  139. try {
  140. const [url] = await this.storage
  141. .bucket(bucketName)
  142. .file(key)
  143. .getSignedUrl({
  144. action: 'read',
  145. expires: Date.now() + this.settings.signedUrlExpiryInMs,
  146. })
  147. return url
  148. } catch (err) {
  149. throw PersistorHelper.wrapError(
  150. err,
  151. 'error generating signed url for GCS file',
  152. { bucketName, key },
  153. ReadError
  154. )
  155. }
  156. }
  157. async getObjectSize(bucketName, key) {
  158. try {
  159. const [metadata] = await this.storage
  160. .bucket(bucketName)
  161. .file(key)
  162. .getMetadata()
  163. return metadata.size
  164. } catch (err) {
  165. throw PersistorHelper.wrapError(
  166. err,
  167. 'error getting size of GCS object',
  168. { bucketName, key },
  169. ReadError
  170. )
  171. }
  172. }
  173. async getObjectMd5Hash(bucketName, key) {
  174. try {
  175. const [metadata] = await this.storage
  176. .bucket(bucketName)
  177. .file(key)
  178. .getMetadata()
  179. return PersistorHelper.base64ToHex(metadata.md5Hash)
  180. } catch (err) {
  181. throw PersistorHelper.wrapError(
  182. err,
  183. 'error getting hash of GCS object',
  184. { bucketName, key },
  185. ReadError
  186. )
  187. }
  188. }
  189. async deleteObject(bucketName, key) {
  190. try {
  191. const file = this.storage.bucket(bucketName).file(key)
  192. if (this.settings.deletedBucketSuffix) {
  193. await file.copy(
  194. this.storage
  195. .bucket(`${bucketName}${this.settings.deletedBucketSuffix}`)
  196. .file(`${key}-${new Date().toISOString()}`)
  197. )
  198. }
  199. if (this.settings.unlockBeforeDelete) {
  200. await file.setMetadata({ eventBasedHold: false })
  201. }
  202. await file.delete()
  203. } catch (err) {
  204. // ignore 404s: it's fine if the file doesn't exist.
  205. if (err.code === 404) {
  206. return
  207. }
  208. throw PersistorHelper.wrapError(
  209. err,
  210. 'error deleting GCS object',
  211. { bucketName, key },
  212. WriteError
  213. )
  214. }
  215. }
  216. async deleteDirectory(bucketName, key) {
  217. const prefix = ensurePrefixIsDirectory(key)
  218. let query = { prefix, autoPaginate: false }
  219. do {
  220. try {
  221. const [files, nextQuery] = await this.storage
  222. .bucket(bucketName)
  223. .getFiles(query)
  224. // iterate over paginated results using the nextQuery returned by getFiles
  225. query = nextQuery
  226. if (Array.isArray(files) && files.length > 0) {
  227. await asyncPool(
  228. this.settings.deleteConcurrency,
  229. files,
  230. async file => {
  231. await this.deleteObject(bucketName, file.name)
  232. }
  233. )
  234. }
  235. } catch (err) {
  236. const error = PersistorHelper.wrapError(
  237. err,
  238. 'failed to delete directory in GCS',
  239. { bucketName, key },
  240. WriteError
  241. )
  242. if (error instanceof NotFoundError) {
  243. return
  244. }
  245. throw error
  246. }
  247. } while (query)
  248. }
  249. async directorySize(bucketName, key) {
  250. let files
  251. const prefix = ensurePrefixIsDirectory(key)
  252. try {
  253. const [response] = await this.storage
  254. .bucket(bucketName)
  255. .getFiles({ prefix })
  256. files = response
  257. } catch (err) {
  258. throw PersistorHelper.wrapError(
  259. err,
  260. 'failed to list objects in GCS',
  261. { bucketName, key },
  262. ReadError
  263. )
  264. }
  265. return files.reduce((acc, file) => Number(file.metadata.size) + acc, 0)
  266. }
  267. async checkIfObjectExists(bucketName, key) {
  268. try {
  269. const [response] = await this.storage
  270. .bucket(bucketName)
  271. .file(key)
  272. .exists()
  273. return response
  274. } catch (err) {
  275. throw PersistorHelper.wrapError(
  276. err,
  277. 'error checking if file exists in GCS',
  278. { bucketName, key },
  279. ReadError
  280. )
  281. }
  282. }
  283. async copyObject(bucketName, sourceKey, destKey) {
  284. try {
  285. const src = this.storage.bucket(bucketName).file(sourceKey)
  286. const dest = this.storage.bucket(bucketName).file(destKey)
  287. await src.copy(dest)
  288. } catch (err) {
  289. // fake-gcs-server has a bug that returns an invalid response when the file does not exist
  290. if (err.message === 'Cannot parse response as JSON: not found\n') {
  291. err.code = 404
  292. }
  293. throw PersistorHelper.wrapError(
  294. err,
  295. 'failed to copy file in GCS',
  296. { bucketName, sourceKey, destKey },
  297. WriteError
  298. )
  299. }
  300. }
  301. }
  302. function ensurePrefixIsDirectory(key) {
  303. return key === '' || key.endsWith('/') ? key : `${key}/`
  304. }