Feat/record with pipeline (#318)

* use pipeline for nodejs streams

* use pipeline for nodejs streams
This commit is contained in:
Hoan Luu Huu
2024-04-30 07:39:24 -04:00
committed by GitHub
parent b765232d4f
commit 3b47162d13
+12 -14
View File
@@ -3,6 +3,7 @@ const Websocket = require('ws');
const PCMToMP3Encoder = require('./encoder'); const PCMToMP3Encoder = require('./encoder');
const wav = require('wav'); const wav = require('wav');
const { getUploader } = require('./utils'); const { getUploader } = require('./utils');
const { pipeline } = require('stream');
async function upload(logger, socket) { async function upload(logger, socket) {
socket._recvInitialMetadata = false; socket._recvInitialMetadata = false;
@@ -60,22 +61,19 @@ async function upload(logger, socket) {
bitrate: 128 bitrate: 128
}, logger); }, logger);
} }
const handleError = (err, streamType) => {
logger.error(
{ err },
`Error while streaming for vendor: ${obj.vendor}, pipe: ${streamType}: ${err.message}`
);
};
/* start streaming data */ /* start streaming data */
const duplex = Websocket.createWebSocketStream(socket); pipeline(
duplex Websocket.createWebSocketStream(socket),
.on('error', (err) => handleError(err, 'duplex')) encoder,
.pipe(encoder) uploadStream,
.on('error', (err) => handleError(err, 'encoder')) (error) => {
.pipe(uploadStream) if (error) {
.on('error', (err) => handleError(err, 'uploadStream')); logger.error({ error }, 'pipeline error, cannot upload data to storage');
socket.close();
}
}
);
} else { } else {
logger.info(`account ${accountSid} does not have any bucket credential, close the socket`); logger.info(`account ${accountSid} does not have any bucket credential, close the socket`);
socket.close(); socket.close();