strong-tie/inbound-calls
0
1'use strict'2 3const { test } = require('tap')4const fs = require('fs')5const proxyquire = require('proxyquire')6const { file, runTests } = require('./helper')7 8const MAX_WRITE = 16 * 10249 10runTests(buildTests)11 12function buildTests (test, sync) {13 // Reset the umask for testing14 process.umask(0o000)15 test('retry on EAGAIN', (t) => {16 t.plan(7)17 18 const fakeFs = Object.create(fs)19 fakeFs.write = function (fd, buf, ...args) {20 t.pass('fake fs.write called')21 fakeFs.write = fs.write22 const err = new Error('EAGAIN')23 err.code = 'EAGAIN'24 process.nextTick(args.pop(), err)25 }26 const SonicBoom = proxyquire('../', {27 fs: fakeFs28 })29 30 const dest = file()31 const fd = fs.openSync(dest, 'w')32 const stream = new SonicBoom({ fd, sync: false, minLength: 0 })33 34 stream.on('ready', () => {35 t.pass('ready emitted')36 })37 38 t.ok(stream.write('hello world\n'))39 t.ok(stream.write('something else\n'))40 41 stream.end()42 43 stream.on('finish', () => {44 fs.readFile(dest, 'utf8', (err, data) => {45 t.error(err)46 t.equal(data, 'hello world\nsomething else\n')47 })48 })49 stream.on('close', () => {50 t.pass('close emitted')51 })52 })53}54 55test('emit error on async EAGAIN', (t) => {56 t.plan(11)57 58 const fakeFs = Object.create(fs)59 fakeFs.write = function (fd, buf, ...args) {60 t.pass('fake fs.write called')61 fakeFs.write = fs.write62 const err = new Error('EAGAIN')63 err.code = 'EAGAIN'64 process.nextTick(args[args.length - 1], err)65 }66 const SonicBoom = proxyquire('../', {67 fs: fakeFs68 })69 70 const dest = file()71 const fd = fs.openSync(dest, 'w')72 const stream = new SonicBoom({73 fd,74 sync: false,75 minLength: 12,76 retryEAGAIN: (err, writeBufferLen, remainingBufferLen) => {77 t.equal(err.code, 'EAGAIN')78 t.equal(writeBufferLen, 12)79 t.equal(remainingBufferLen, 0)80 return false81 }82 })83 84 stream.on('ready', () => {85 t.pass('ready emitted')86 })87 88 stream.once('error', err => {89 t.equal(err.code, 'EAGAIN')90 t.ok(stream.write('something else\n'))91 })92 93 t.ok(stream.write('hello world\n'))94 95 stream.end()96 97 stream.on('finish', () => {98 fs.readFile(dest, 'utf8', (err, data) => {99 t.error(err)100 t.equal(data, 'hello world\nsomething else\n')101 })102 })103 stream.on('close', () => {104 t.pass('close emitted')105 })106})107 108test('retry on EAGAIN (sync)', (t) => {109 t.plan(7)110 111 const fakeFs = Object.create(fs)112 fakeFs.writeSync = function (fd, buf, enc) {113 t.pass('fake fs.writeSync called')114 fakeFs.writeSync = fs.writeSync115 const err = new Error('EAGAIN')116 err.code = 'EAGAIN'117 throw err118 }119 const SonicBoom = proxyquire('../', {120 fs: fakeFs121 })122 123 const dest = file()124 const fd = fs.openSync(dest, 'w')125 const stream = new SonicBoom({ fd, minLength: 0, sync: true })126 127 stream.on('ready', () => {128 t.pass('ready emitted')129 })130 131 t.ok(stream.write('hello world\n'))132 t.ok(stream.write('something else\n'))133 134 stream.end()135 136 stream.on('finish', () => {137 fs.readFile(dest, 'utf8', (err, data) => {138 t.error(err)139 t.equal(data, 'hello world\nsomething else\n')140 })141 })142 stream.on('close', () => {143 t.pass('close emitted')144 })145})146 147test('emit error on EAGAIN (sync)', (t) => {148 t.plan(11)149 150 const fakeFs = Object.create(fs)151 fakeFs.writeSync = function (fd, buf, enc) {152 t.pass('fake fs.writeSync called')153 fakeFs.writeSync = fs.writeSync154 const err = new Error('EAGAIN')155 err.code = 'EAGAIN'156 throw err157 }158 const SonicBoom = proxyquire('../', {159 fs: fakeFs160 })161 162 const dest = file()163 const fd = fs.openSync(dest, 'w')164 const stream = new SonicBoom({165 fd,166 minLength: 0,167 sync: true,168 retryEAGAIN: (err, writeBufferLen, remainingBufferLen) => {169 t.equal(err.code, 'EAGAIN')170 t.equal(writeBufferLen, 12)171 t.equal(remainingBufferLen, 0)172 return false173 }174 })175 176 stream.on('ready', () => {177 t.pass('ready emitted')178 })179 180 stream.once('error', err => {181 t.equal(err.code, 'EAGAIN')182 t.ok(stream.write('something else\n'))183 })184 185 t.ok(stream.write('hello world\n'))186 187 stream.end()188 189 stream.on('finish', () => {190 fs.readFile(dest, 'utf8', (err, data) => {191 t.error(err)192 t.equal(data, 'hello world\nsomething else\n')193 })194 })195 stream.on('close', () => {196 t.pass('close emitted')197 })198})199 200test('retryEAGAIN receives remaining buffer on async if write fails', (t) => {201 t.plan(12)202 203 const fakeFs = Object.create(fs)204 const SonicBoom = proxyquire('../', {205 fs: fakeFs206 })207 208 const dest = file()209 const fd = fs.openSync(dest, 'w')210 const stream = new SonicBoom({211 fd,212 sync: false,213 minLength: 12,214 retryEAGAIN: (err, writeBufferLen, remainingBufferLen) => {215 t.equal(err.code, 'EAGAIN')216 t.equal(writeBufferLen, 12)217 t.equal(remainingBufferLen, 11)218 return false219 }220 })221 222 stream.on('ready', () => {223 t.pass('ready emitted')224 })225 226 stream.once('error', err => {227 t.equal(err.code, 'EAGAIN')228 t.ok(stream.write('done'))229 })230 231 fakeFs.write = function (fd, buf, ...args) {232 t.pass('fake fs.write called')233 fakeFs.write = fs.write234 const err = new Error('EAGAIN')235 err.code = 'EAGAIN'236 t.ok(stream.write('sonic boom\n'))237 process.nextTick(args[args.length - 1], err)238 }239 240 t.ok(stream.write('hello world\n'))241 242 stream.end()243 244 stream.on('finish', () => {245 fs.readFile(dest, 'utf8', (err, data) => {246 t.error(err)247 t.equal(data, 'hello world\nsonic boom\ndone')248 })249 })250 stream.on('close', () => {251 t.pass('close emitted')252 })253})254 255test('retryEAGAIN receives remaining buffer if exceeds maxWrite', (t) => {256 t.plan(17)257 258 const fakeFs = Object.create(fs)259 const SonicBoom = proxyquire('../', {260 fs: fakeFs261 })262 263 const dest = file()264 const fd = fs.openSync(dest, 'w')265 const buf = Buffer.alloc(MAX_WRITE - 2).fill('x').toString() // 1 MB266 const stream = new SonicBoom({267 fd,268 sync: false,269 minLength: MAX_WRITE - 1,270 retryEAGAIN: (err, writeBufferLen, remainingBufferLen) => {271 t.equal(err.code, 'EAGAIN', 'retryEAGAIN received EAGAIN error')272 t.equal(writeBufferLen, buf.length, 'writeBufferLen === buf.length')273 t.equal(remainingBufferLen, 23, 'remainingBufferLen === 23')274 return false275 }276 })277 278 stream.on('ready', () => {279 t.pass('ready emitted')280 })281 282 fakeFs.write = function (fd, buf, ...args) {283 t.pass('fake fs.write called')284 const err = new Error('EAGAIN')285 err.code = 'EAGAIN'286 process.nextTick(args.pop(), err)287 }288 289 fakeFs.writeSync = function (fd, buf, enc) {290 t.pass('fake fs.write called')291 const err = new Error('EAGAIN')292 err.code = 'EAGAIN'293 throw err294 }295 296 t.ok(stream.write(buf), 'write buf')297 t.notOk(stream.write('hello world\nsonic boom\n'), 'write hello world sonic boom')298 299 stream.once('error', err => {300 t.equal(err.code, 'EAGAIN', 'bubbled error should be EAGAIN')301 302 try {303 stream.flushSync()304 } catch (err) {305 t.equal(err.code, 'EAGAIN', 'thrown error should be EAGAIN')306 fakeFs.write = fs.write307 fakeFs.writeSync = fs.writeSync308 stream.end()309 }310 })311 312 stream.on('finish', () => {313 t.pass('finish emitted')314 fs.readFile(dest, 'utf8', (err, data) => {315 t.error(err)316 t.equal(data, `${buf}hello world\nsonic boom\n`, 'data on file should match written')317 })318 })319 stream.on('close', () => {320 t.pass('close emitted')321 })322})323 324test('retry on EBUSY', (t) => {325 t.plan(7)326 327 const fakeFs = Object.create(fs)328 fakeFs.write = function (fd, buf, ...args) {329 t.pass('fake fs.write called')330 fakeFs.write = fs.write331 const err = new Error('EBUSY')332 err.code = 'EBUSY'333 process.nextTick(args.pop(), err)334 }335 const SonicBoom = proxyquire('..', {336 fs: fakeFs337 })338 339 const dest = file()340 const fd = fs.openSync(dest, 'w')341 const stream = new SonicBoom({ fd, sync: false, minLength: 0 })342 343 stream.on('ready', () => {344 t.pass('ready emitted')345 })346 347 t.ok(stream.write('hello world\n'))348 t.ok(stream.write('something else\n'))349 350 stream.end()351 352 stream.on('finish', () => {353 fs.readFile(dest, 'utf8', (err, data) => {354 t.error(err)355 t.equal(data, 'hello world\nsomething else\n')356 })357 })358 stream.on('close', () => {359 t.pass('close emitted')360 })361})362 363test('emit error on async EBUSY', (t) => {364 t.plan(11)365 366 const fakeFs = Object.create(fs)367 fakeFs.write = function (fd, buf, ...args) {368 t.pass('fake fs.write called')369 fakeFs.write = fs.write370 const err = new Error('EBUSY')371 err.code = 'EBUSY'372 process.nextTick(args.pop(), err)373 }374 const SonicBoom = proxyquire('..', {375 fs: fakeFs376 })377 378 const dest = file()379 const fd = fs.openSync(dest, 'w')380 const stream = new SonicBoom({381 fd,382 sync: false,383 minLength: 12,384 retryEAGAIN: (err, writeBufferLen, remainingBufferLen) => {385 t.equal(err.code, 'EBUSY')386 t.equal(writeBufferLen, 12)387 t.equal(remainingBufferLen, 0)388 return false389 }390 })391 392 stream.on('ready', () => {393 t.pass('ready emitted')394 })395 396 stream.once('error', err => {397 t.equal(err.code, 'EBUSY')398 t.ok(stream.write('something else\n'))399 })400 401 t.ok(stream.write('hello world\n'))402 403 stream.end()404 405 stream.on('finish', () => {406 fs.readFile(dest, 'utf8', (err, data) => {407 t.error(err)408 t.equal(data, 'hello world\nsomething else\n')409 })410 })411 stream.on('close', () => {412 t.pass('close emitted')413 })414})415 