GcsPersistor.js 7.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287
  1. const fs = require('fs')
  2. const { promisify } = require('util')
  3. const Stream = require('stream')
  4. const { Storage } = 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 pipeline = promisify(Stream.pipeline)
  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. this.storage = new Storage(this.settings.endpoint)
  17. // workaround for broken uploads with custom endpoints:
  18. // https://github.com/googleapis/nodejs-storage/issues/898
  19. if (this.settings.endpoint && this.settings.endpoint.apiEndpoint) {
  20. this.storage.interceptors.push({
  21. request: (reqOpts) => {
  22. const url = new URL(reqOpts.uri)
  23. url.host = this.settings.endpoint.apiEndpoint
  24. if (this.settings.endpoint.apiScheme) {
  25. url.protocol = this.settings.endpoint.apiScheme
  26. }
  27. reqOpts.uri = url.toString()
  28. return reqOpts
  29. }
  30. })
  31. }
  32. }
  33. async sendFile(bucketName, key, fsPath) {
  34. return this.sendStream(bucketName, key, fs.createReadStream(fsPath))
  35. }
  36. async sendStream(bucketName, key, readStream, opts = {}) {
  37. try {
  38. // egress from us to gcs
  39. const observeOptions = {
  40. metric: 'gcs.egress',
  41. Metrics: this.settings.Metrics
  42. }
  43. let sourceMd5 = opts.sourceMd5
  44. if (!sourceMd5) {
  45. // if there is no supplied md5 hash, we calculate the hash as the data passes through
  46. observeOptions.hash = 'md5'
  47. }
  48. const observer = new PersistorHelper.ObserverStream(observeOptions)
  49. const writeOptions = {
  50. // disabling of resumable uploads is recommended by Google:
  51. resumable: false
  52. }
  53. if (sourceMd5) {
  54. writeOptions.validation = 'md5'
  55. writeOptions.metadata = writeOptions.metadata || {}
  56. writeOptions.metadata.md5Hash = PersistorHelper.hexToBase64(sourceMd5)
  57. }
  58. if (opts.contentType) {
  59. writeOptions.metadata = writeOptions.metadata || {}
  60. writeOptions.metadata.contentType = opts.contentType
  61. }
  62. if (opts.contentEncoding) {
  63. writeOptions.metadata = writeOptions.metadata || {}
  64. writeOptions.metadata.contentEncoding = opts.contentEncoding
  65. }
  66. const uploadStream = this.storage
  67. .bucket(bucketName)
  68. .file(key)
  69. .createWriteStream(writeOptions)
  70. await pipeline(readStream, observer, uploadStream)
  71. // if we didn't have an md5 hash, we should compare our computed one with Google's
  72. // as we couldn't tell GCS about it beforehand
  73. if (!sourceMd5) {
  74. sourceMd5 = observer.getHash()
  75. // throws on mismatch
  76. await PersistorHelper.verifyMd5(this, bucketName, key, sourceMd5)
  77. }
  78. } catch (err) {
  79. throw PersistorHelper.wrapError(
  80. err,
  81. 'upload to GCS failed',
  82. { bucketName, key },
  83. WriteError
  84. )
  85. }
  86. }
  87. async getObjectStream(bucketName, key, opts = {}) {
  88. const stream = this.storage
  89. .bucket(bucketName)
  90. .file(key)
  91. .createReadStream(opts)
  92. // ingress to us from gcs
  93. const observer = new PersistorHelper.ObserverStream({
  94. metric: 'gcs.ingress',
  95. Metrics: this.settings.Metrics
  96. })
  97. try {
  98. // wait for the pipeline to be ready, to catch non-200s
  99. await PersistorHelper.getReadyPipeline(stream, observer)
  100. return observer
  101. } catch (err) {
  102. throw PersistorHelper.wrapError(
  103. err,
  104. 'error reading file from GCS',
  105. { bucketName, key, opts },
  106. ReadError
  107. )
  108. }
  109. }
  110. async getRedirectUrl(bucketName, key) {
  111. try {
  112. const [url] = await this.storage
  113. .bucket(bucketName)
  114. .file(key)
  115. .getSignedUrl({
  116. action: 'read',
  117. expires: Date.now() + this.settings.signedUrlExpiryInMs
  118. })
  119. return url
  120. } catch (err) {
  121. throw PersistorHelper.wrapError(
  122. err,
  123. 'error generating signed url for GCS file',
  124. { bucketName, key },
  125. ReadError
  126. )
  127. }
  128. }
  129. async getObjectSize(bucketName, key) {
  130. try {
  131. const [metadata] = await this.storage
  132. .bucket(bucketName)
  133. .file(key)
  134. .getMetadata()
  135. return metadata.size
  136. } catch (err) {
  137. throw PersistorHelper.wrapError(
  138. err,
  139. 'error getting size of GCS object',
  140. { bucketName, key },
  141. ReadError
  142. )
  143. }
  144. }
  145. async getObjectMd5Hash(bucketName, key) {
  146. try {
  147. const [metadata] = await this.storage
  148. .bucket(bucketName)
  149. .file(key)
  150. .getMetadata()
  151. return PersistorHelper.base64ToHex(metadata.md5Hash)
  152. } catch (err) {
  153. throw PersistorHelper.wrapError(
  154. err,
  155. 'error getting hash of GCS object',
  156. { bucketName, key },
  157. ReadError
  158. )
  159. }
  160. }
  161. async deleteObject(bucketName, key) {
  162. try {
  163. const file = this.storage.bucket(bucketName).file(key)
  164. if (this.settings.deletedBucketSuffix) {
  165. await file.copy(
  166. this.storage
  167. .bucket(`${bucketName}${this.settings.deletedBucketSuffix}`)
  168. .file(`${key}-${new Date().toISOString()}`)
  169. )
  170. }
  171. if (this.settings.unlockBeforeDelete) {
  172. await file.setMetadata({ eventBasedHold: false })
  173. }
  174. await file.delete()
  175. } catch (err) {
  176. const error = PersistorHelper.wrapError(
  177. err,
  178. 'error deleting GCS object',
  179. { bucketName, key },
  180. WriteError
  181. )
  182. if (!(error instanceof NotFoundError)) {
  183. throw error
  184. }
  185. }
  186. }
  187. async deleteDirectory(bucketName, key) {
  188. try {
  189. const [files] = await this.storage
  190. .bucket(bucketName)
  191. .getFiles({ directory: key })
  192. await asyncPool(this.settings.deleteConcurrency, files, async (file) => {
  193. await this.deleteObject(bucketName, file.name)
  194. })
  195. } catch (err) {
  196. const error = PersistorHelper.wrapError(
  197. err,
  198. 'failed to delete directory in GCS',
  199. { bucketName, key },
  200. WriteError
  201. )
  202. if (error instanceof NotFoundError) {
  203. return
  204. }
  205. throw error
  206. }
  207. }
  208. async directorySize(bucketName, key) {
  209. let files
  210. try {
  211. const [response] = await this.storage
  212. .bucket(bucketName)
  213. .getFiles({ directory: key })
  214. files = response
  215. } catch (err) {
  216. throw PersistorHelper.wrapError(
  217. err,
  218. 'failed to list objects in GCS',
  219. { bucketName, key },
  220. ReadError
  221. )
  222. }
  223. return files.reduce((acc, file) => Number(file.metadata.size) + acc, 0)
  224. }
  225. async checkIfObjectExists(bucketName, key) {
  226. try {
  227. const [response] = await this.storage
  228. .bucket(bucketName)
  229. .file(key)
  230. .exists()
  231. return response
  232. } catch (err) {
  233. throw PersistorHelper.wrapError(
  234. err,
  235. 'error checking if file exists in GCS',
  236. { bucketName, key },
  237. ReadError
  238. )
  239. }
  240. }
  241. async copyObject(bucketName, sourceKey, destKey) {
  242. try {
  243. const src = this.storage.bucket(bucketName).file(sourceKey)
  244. const dest = this.storage.bucket(bucketName).file(destKey)
  245. await src.copy(dest)
  246. } catch (err) {
  247. // fake-gcs-server has a bug that returns an invalid response when the file does not exist
  248. if (err.message === 'Cannot parse response as JSON: not found\n') {
  249. err.code = 404
  250. }
  251. throw PersistorHelper.wrapError(
  252. err,
  253. 'failed to copy file in GCS',
  254. { bucketName, sourceKey, destKey },
  255. WriteError
  256. )
  257. }
  258. }
  259. }