How to keep the request open to use the write() method after a long time

Viewed 572

I need to keep the connection open so after I finish the music I write the new data. The problem is that the way I did, the stream simply stops after the first song. How can I keep the connection open and play the next songs too?

const fs = require('fs');
const express = require('express');
const app = express();
const server = require('http').createServer(app)
const getMP3Duration = require('get-mp3-duration')

let sounds = ['61880.mp3', '62026.mp3', '62041.mp3', '62090.mp3', '62257.mp3', '60763.mp3']

app.get('/current', async (req, res) => {
    let readStream = fs.createReadStream('sounds/61068.mp3')
    let duration = await getMP3Duration(fs.readFileSync('sounds/61068.mp3'))

    let pipe = readStream.pipe(res, {end: false})

    async function put(){
        let file_path = 'sounds/'+sounds[Math.random() * sounds.length-1]

        duration = await getMP3Duration(fs.readFileSync(file_path))

        readStream = fs.createReadStream(file_path)

        readStream.on('data', chunk => {
            console.log(chunk)
            pipe.write(chunk)
        })

        console.log('Current Sound: ', file_path)

        setTimeout(put, duration)
    }

    setTimeout(put, duration)
})

server.listen(3005, async function () {
    console.log('Server is running on port 3005...')
});
3 Answers

You should use a library or look at the source code and see what they do. A good one is: https://github.com/obastemur/mediaserver

TIP:
Always start your research by learning from other projects.. (When possible or when you are not inventing the wheel ;)) you are not the first to do so or to hit this problem :)

a quick search with the phrase "nodejs stream mp3 github" gave me few directions.. Good luck !

Express works by returning a single response to a single request. As soon as the request has been sent, a new request needs to be generated to trigger a new response.

In your case however you want to keep on generating new responses out of a single request.

Two approaches can be used to solve your problem:

  1. Change the way you create your response to satisfy your use-case.
  2. use an instantaneous communication framework (websocket). The best and simplest which comes to my mind is socket.io

Adapting express

The solution here is to follow this procedure:

  1. Request on endpoint /current comes in
  2. The audio sequence is prepared
  3. The stream of the entire sequence is returned

So your handler would look like that:

const fs = require('fs');
const express = require('express');
const app = express();
const server = require('http').createServer(app);
// Import the PassThrough class to concatenate the streams
const { PassThrough } = require('stream');
// The array of sounds now contain all the sounds
const sounds = ['61068.mp3','61880.mp3', '62026.mp3', '62041.mp3', '62090.mp3', '62257.mp3', '60763.mp3'];


// function which concatenate an array of streams
const concatStreams = streamArray => {
  let pass = new PassThrough();
  let waiting = streamArray.length;
  streamArray.forEach(soundStream => {
    pass = soundStream.pipe(pass, {end: false});
    soundStream.once('end', () => --waiting === 0 && pass.emit('end'));
  });
  return pass;
};

// function which returns a shuffled array
const shuffle = (array) => {
  const a = [...array]; // shallow copy of the array
  for (let i = a.length - 1; i > 0; i--) {
    const j = Math.floor(Math.random() * (i + 1));
    [a[i], a[j]] = [a[j], a[i]];
  }
  return a;
};


server.get('/current', (req, res) => {
  // Start by shuffling the array
  const shuffledSounds = shuffle(sounds);

  // Create a readable stream for each sound
  const streams = shuffledSounds.map(sound => fs.createReadStream(`sounds/${sound}`));

  // Concatenate all the streams into a single stream
  const readStream = concatStreams(streams);

  // This will wait until we know the readable stream is actually valid before piping
  readStream.on('open', function () {
    // This just pipes the read stream to the response object (which goes to the client)
    // the response is automatically ended when the stream emits the "end" event
    readStream.pipe(res);
  });
});

Notice that the function does not require the async keyword any longer. The process is still asynchronous but the coding is emitter based instead of promise based.

If you want to loop the sounds you can create additional steps of shuffling/mapping to stream/concatenation.

I did not include the socketio alternative as to keep it simple.

Final Solution After a Few Edits:

I suspect your main issue is with your random array element generator. You need to wrap what you have with Math.floor to round down to ensure you end up with a whole number:

sounds[Math.floor(Math.random() * sounds.length)]

Also, Readstream.pipe returns the destination, so what you're doing makes sense. However, you might get unexpected results with calling on('data') on your readable after you've already piped from it. The node.js streams docs mention this. I tested out your code on my local machine and it doesn't seem to be an issue, but it might make sense to change this so you don't have problems in the future.

Choose One API Style

The Readable stream API evolved across multiple Node.js versions and provides multiple methods of consuming stream data. In general, developers should choose one of the methods of consuming data and should never use multiple methods to consume data from a single stream. Specifically, using a combination of on('data'), on('readable'), pipe(), or async iterators could lead to unintuitive behavior.

Instead of calling on('data') and res.write, I would just pipe from the readStream into the res again. Also, unless you really want to get the duration, I would pull that library out and just use the readStream.end event to make additional calls to put(). This works because you're passing the false option when piping, which disables the default end event functionality on the write stream and leaves it open. However, it still gets emitted, so you can use that as a marker to know when the readable has finished piping. Here's the refactored code:

const fs = require('fs');
const express = require('express');
const app = express();
const server = require('http').createServer(app)
//const getMP3Duration = require('get-mp3-duration') no longer needed

let sounds = ['61880.mp3', '62026.mp3', '62041.mp3', '62090.mp3', '62257.mp3', '60763.mp3']

app.get('/current', async (req, res) => {
    let readStream = fs.createReadStream('sounds/61068.mp3')
    let duration = await getMP3Duration(fs.readFileSync('sounds/61068.mp3'))

    let pipe = readStream.pipe(res, {end: false})

    function put(){
        let file_path = 'sounds/'+sounds[Math.floor(Math.random() * sounds.length)]

        readStream = fs.createReadStream(file_path)
        
        // you may also be able to do readStream.pipe(res, {end: false})
        readStream.pipe(pipe, {end: false})

        console.log('Current Sound: ', file_path)

        readStream.on('end', () => {
            put()
        });
    }

    readStream.on('end', () => {
        put()
    });
})

server.listen(3005, async function () {
    console.log('Server is running on port 3005...')
});
Related