Blame view

node_modules/graceful-fs/graceful-fs.js 12.4 KB
7820380e   “wangming”   1
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
  var fs = require('fs')
  var polyfills = require('./polyfills.js')
  var legacy = require('./legacy-streams.js')
  var clone = require('./clone.js')
  
  var util = require('util')
  
  /* istanbul ignore next - node 0.x polyfill */
  var gracefulQueue
  var previousSymbol
  
  /* istanbul ignore else - node 0.x polyfill */
  if (typeof Symbol === 'function' && typeof Symbol.for === 'function') {
    gracefulQueue = Symbol.for('graceful-fs.queue')
    // This is used in testing by future versions
    previousSymbol = Symbol.for('graceful-fs.previous')
  } else {
    gracefulQueue = '___graceful-fs.queue'
    previousSymbol = '___graceful-fs.previous'
  }
  
  function noop () {}
  
  function publishQueue(context, queue) {
    Object.defineProperty(context, gracefulQueue, {
      get: function() {
        return queue
      }
    })
  }
  
  var debug = noop
  if (util.debuglog)
    debug = util.debuglog('gfs4')
  else if (/\bgfs4\b/i.test(process.env.NODE_DEBUG || ''))
    debug = function() {
      var m = util.format.apply(util, arguments)
      m = 'GFS4: ' + m.split(/\n/).join('\nGFS4: ')
      console.error(m)
    }
  
  // Once time initialization
  if (!fs[gracefulQueue]) {
    // This queue can be shared by multiple loaded instances
    var queue = global[gracefulQueue] || []
    publishQueue(fs, queue)
  
    // Patch fs.close/closeSync to shared queue version, because we need
    // to retry() whenever a close happens *anywhere* in the program.
    // This is essential when multiple graceful-fs instances are
    // in play at the same time.
    fs.close = (function (fs$close) {
      function close (fd, cb) {
        return fs$close.call(fs, fd, function (err) {
          // This function uses the graceful-fs shared queue
          if (!err) {
            resetQueue()
          }
  
          if (typeof cb === 'function')
            cb.apply(this, arguments)
        })
      }
  
      Object.defineProperty(close, previousSymbol, {
        value: fs$close
      })
      return close
    })(fs.close)
  
    fs.closeSync = (function (fs$closeSync) {
      function closeSync (fd) {
        // This function uses the graceful-fs shared queue
        fs$closeSync.apply(fs, arguments)
        resetQueue()
      }
  
      Object.defineProperty(closeSync, previousSymbol, {
        value: fs$closeSync
      })
      return closeSync
    })(fs.closeSync)
  
    if (/\bgfs4\b/i.test(process.env.NODE_DEBUG || '')) {
      process.on('exit', function() {
        debug(fs[gracefulQueue])
        require('assert').equal(fs[gracefulQueue].length, 0)
      })
    }
  }
  
  if (!global[gracefulQueue]) {
    publishQueue(global, fs[gracefulQueue]);
  }
  
  module.exports = patch(clone(fs))
  if (process.env.TEST_GRACEFUL_FS_GLOBAL_PATCH && !fs.__patched) {
      module.exports = patch(fs)
      fs.__patched = true;
  }
  
  function patch (fs) {
    // Everything that references the open() function needs to be in here
    polyfills(fs)
    fs.gracefulify = patch
  
    fs.createReadStream = createReadStream
    fs.createWriteStream = createWriteStream
    var fs$readFile = fs.readFile
    fs.readFile = readFile
    function readFile (path, options, cb) {
      if (typeof options === 'function')
        cb = options, options = null
  
      return go$readFile(path, options, cb)
  
      function go$readFile (path, options, cb, startTime) {
        return fs$readFile(path, options, function (err) {
          if (err && (err.code === 'EMFILE' || err.code === 'ENFILE'))
            enqueue([go$readFile, [path, options, cb], err, startTime || Date.now(), Date.now()])
          else {
            if (typeof cb === 'function')
              cb.apply(this, arguments)
          }
        })
      }
    }
  
    var fs$writeFile = fs.writeFile
    fs.writeFile = writeFile
    function writeFile (path, data, options, cb) {
      if (typeof options === 'function')
        cb = options, options = null
  
      return go$writeFile(path, data, options, cb)
  
      function go$writeFile (path, data, options, cb, startTime) {
        return fs$writeFile(path, data, options, function (err) {
          if (err && (err.code === 'EMFILE' || err.code === 'ENFILE'))
            enqueue([go$writeFile, [path, data, options, cb], err, startTime || Date.now(), Date.now()])
          else {
            if (typeof cb === 'function')
              cb.apply(this, arguments)
          }
        })
      }
    }
  
    var fs$appendFile = fs.appendFile
    if (fs$appendFile)
      fs.appendFile = appendFile
    function appendFile (path, data, options, cb) {
      if (typeof options === 'function')
        cb = options, options = null
  
      return go$appendFile(path, data, options, cb)
  
      function go$appendFile (path, data, options, cb, startTime) {
        return fs$appendFile(path, data, options, function (err) {
          if (err && (err.code === 'EMFILE' || err.code === 'ENFILE'))
            enqueue([go$appendFile, [path, data, options, cb], err, startTime || Date.now(), Date.now()])
          else {
            if (typeof cb === 'function')
              cb.apply(this, arguments)
          }
        })
      }
    }
  
    var fs$copyFile = fs.copyFile
    if (fs$copyFile)
      fs.copyFile = copyFile
    function copyFile (src, dest, flags, cb) {
      if (typeof flags === 'function') {
        cb = flags
        flags = 0
      }
      return go$copyFile(src, dest, flags, cb)
  
      function go$copyFile (src, dest, flags, cb, startTime) {
        return fs$copyFile(src, dest, flags, function (err) {
          if (err && (err.code === 'EMFILE' || err.code === 'ENFILE'))
            enqueue([go$copyFile, [src, dest, flags, cb], err, startTime || Date.now(), Date.now()])
          else {
            if (typeof cb === 'function')
              cb.apply(this, arguments)
          }
        })
      }
    }
  
    var fs$readdir = fs.readdir
    fs.readdir = readdir
    var noReaddirOptionVersions = /^v[0-5]\./
    function readdir (path, options, cb) {
      if (typeof options === 'function')
        cb = options, options = null
  
      var go$readdir = noReaddirOptionVersions.test(process.version)
        ? function go$readdir (path, options, cb, startTime) {
          return fs$readdir(path, fs$readdirCallback(
            path, options, cb, startTime
          ))
        }
        : function go$readdir (path, options, cb, startTime) {
          return fs$readdir(path, options, fs$readdirCallback(
            path, options, cb, startTime
          ))
        }
  
      return go$readdir(path, options, cb)
  
      function fs$readdirCallback (path, options, cb, startTime) {
        return function (err, files) {
          if (err && (err.code === 'EMFILE' || err.code === 'ENFILE'))
            enqueue([
              go$readdir,
              [path, options, cb],
              err,
              startTime || Date.now(),
              Date.now()
            ])
          else {
            if (files && files.sort)
              files.sort()
  
            if (typeof cb === 'function')
              cb.call(this, err, files)
          }
        }
      }
    }
  
    if (process.version.substr(0, 4) === 'v0.8') {
      var legStreams = legacy(fs)
      ReadStream = legStreams.ReadStream
      WriteStream = legStreams.WriteStream
    }
  
    var fs$ReadStream = fs.ReadStream
    if (fs$ReadStream) {
      ReadStream.prototype = Object.create(fs$ReadStream.prototype)
      ReadStream.prototype.open = ReadStream$open
    }
  
    var fs$WriteStream = fs.WriteStream
    if (fs$WriteStream) {
      WriteStream.prototype = Object.create(fs$WriteStream.prototype)
      WriteStream.prototype.open = WriteStream$open
    }
  
    Object.defineProperty(fs, 'ReadStream', {
      get: function () {
        return ReadStream
      },
      set: function (val) {
        ReadStream = val
      },
      enumerable: true,
      configurable: true
    })
    Object.defineProperty(fs, 'WriteStream', {
      get: function () {
        return WriteStream
      },
      set: function (val) {
        WriteStream = val
      },
      enumerable: true,
      configurable: true
    })
  
    // legacy names
    var FileReadStream = ReadStream
    Object.defineProperty(fs, 'FileReadStream', {
      get: function () {
        return FileReadStream
      },
      set: function (val) {
        FileReadStream = val
      },
      enumerable: true,
      configurable: true
    })
    var FileWriteStream = WriteStream
    Object.defineProperty(fs, 'FileWriteStream', {
      get: function () {
        return FileWriteStream
      },
      set: function (val) {
        FileWriteStream = val
      },
      enumerable: true,
      configurable: true
    })
  
    function ReadStream (path, options) {
      if (this instanceof ReadStream)
        return fs$ReadStream.apply(this, arguments), this
      else
        return ReadStream.apply(Object.create(ReadStream.prototype), arguments)
    }
  
    function ReadStream$open () {
      var that = this
      open(that.path, that.flags, that.mode, function (err, fd) {
        if (err) {
          if (that.autoClose)
            that.destroy()
  
          that.emit('error', err)
        } else {
          that.fd = fd
          that.emit('open', fd)
          that.read()
        }
      })
    }
  
    function WriteStream (path, options) {
      if (this instanceof WriteStream)
        return fs$WriteStream.apply(this, arguments), this
      else
        return WriteStream.apply(Object.create(WriteStream.prototype), arguments)
    }
  
    function WriteStream$open () {
      var that = this
      open(that.path, that.flags, that.mode, function (err, fd) {
        if (err) {
          that.destroy()
          that.emit('error', err)
        } else {
          that.fd = fd
          that.emit('open', fd)
        }
      })
    }
  
    function createReadStream (path, options) {
      return new fs.ReadStream(path, options)
    }
  
    function createWriteStream (path, options) {
      return new fs.WriteStream(path, options)
    }
  
    var fs$open = fs.open
    fs.open = open
    function open (path, flags, mode, cb) {
      if (typeof mode === 'function')
        cb = mode, mode = null
  
      return go$open(path, flags, mode, cb)
  
      function go$open (path, flags, mode, cb, startTime) {
        return fs$open(path, flags, mode, function (err, fd) {
          if (err && (err.code === 'EMFILE' || err.code === 'ENFILE'))
            enqueue([go$open, [path, flags, mode, cb], err, startTime || Date.now(), Date.now()])
          else {
            if (typeof cb === 'function')
              cb.apply(this, arguments)
          }
        })
      }
    }
  
    return fs
  }
  
  function enqueue (elem) {
    debug('ENQUEUE', elem[0].name, elem[1])
    fs[gracefulQueue].push(elem)
    retry()
  }
  
  // keep track of the timeout between retry() calls
  var retryTimer
  
  // reset the startTime and lastTime to now
  // this resets the start of the 60 second overall timeout as well as the
  // delay between attempts so that we'll retry these jobs sooner
  function resetQueue () {
    var now = Date.now()
    for (var i = 0; i < fs[gracefulQueue].length; ++i) {
      // entries that are only a length of 2 are from an older version, don't
      // bother modifying those since they'll be retried anyway.
      if (fs[gracefulQueue][i].length > 2) {
        fs[gracefulQueue][i][3] = now // startTime
        fs[gracefulQueue][i][4] = now // lastTime
      }
    }
    // call retry to make sure we're actively processing the queue
    retry()
  }
  
  function retry () {
    // clear the timer and remove it to help prevent unintended concurrency
    clearTimeout(retryTimer)
    retryTimer = undefined
  
    if (fs[gracefulQueue].length === 0)
      return
  
    var elem = fs[gracefulQueue].shift()
    var fn = elem[0]
    var args = elem[1]
    // these items may be unset if they were added by an older graceful-fs
    var err = elem[2]
    var startTime = elem[3]
    var lastTime = elem[4]
  
    // if we don't have a startTime we have no way of knowing if we've waited
    // long enough, so go ahead and retry this item now
    if (startTime === undefined) {
      debug('RETRY', fn.name, args)
      fn.apply(null, args)
    } else if (Date.now() - startTime >= 60000) {
      // it's been more than 60 seconds total, bail now
      debug('TIMEOUT', fn.name, args)
      var cb = args.pop()
      if (typeof cb === 'function')
        cb.call(null, err)
    } else {
      // the amount of time between the last attempt and right now
      var sinceAttempt = Date.now() - lastTime
      // the amount of time between when we first tried, and when we last tried
      // rounded up to at least 1
      var sinceStart = Math.max(lastTime - startTime, 1)
      // backoff. wait longer than the total time we've been retrying, but only
      // up to a maximum of 100ms
      var desiredDelay = Math.min(sinceStart * 1.2, 100)
      // it's been long enough since the last retry, do it again
      if (sinceAttempt >= desiredDelay) {
        debug('RETRY', fn.name, args)
        fn.apply(null, args.concat([startTime]))
      } else {
        // if we can't do this job yet, push it to the end of the queue
        // and let the next iteration check again
        fs[gracefulQueue].push(elem)
      }
    }
  
    // schedule our next run if one isn't already scheduled
    if (retryTimer === undefined) {
      retryTimer = setTimeout(retry, 0)
    }
  }