268 lines
8.9 KiB
JavaScript
268 lines
8.9 KiB
JavaScript
const fs = require('fs')
|
|
const exec = require('child_process').exec
|
|
const AWS = require('aws-sdk')
|
|
const Bugsnag = require('@bugsnag/js')
|
|
const uuid = require('uuid').v4
|
|
const version = require('./package.json').version
|
|
const sqs = require('../app/services/sqs')
|
|
const s3 = require('../app/services/s3')
|
|
const userAudioProfileService = require('./user_audio_profile')
|
|
const recordingModel = require('./recording')
|
|
const recordingSalutationModel = require('./recording_salutation')
|
|
const jobService = require('./job')
|
|
const salutationService = require('./salutation')
|
|
let throttleMessageFetching = true
|
|
AWS.config.update({ region: 'us-west-2' })
|
|
const sqsQueueUrl = process.env.SQS_URL
|
|
const mongoUriDev = process.env.MONGODB_URI_DEV
|
|
const mongoUriStaging = process.env.MONGODB_URI_STAGING
|
|
const mongoUriProd = process.env.MONGODB_URI_PROD
|
|
const APP_ENV = process.env.POTION_APP_ENV
|
|
const mongoose = require('mongoose')
|
|
|
|
function execShellCommand(cmd) {
|
|
// const exec = require("child_process").exec;
|
|
return new Promise((resolve, reject) => {
|
|
exec(cmd, { maxBuffer: 1024 * 1000000 }, (error, stdout, stderr) => {
|
|
if (error) {
|
|
console.log('Error while processing python command', error)
|
|
reject(error)
|
|
}
|
|
console.log('Stdout --- ', stdout)
|
|
console.log('Std error --- ', stderr)
|
|
resolve(stdout || stderr)
|
|
})
|
|
})
|
|
}
|
|
|
|
function connectDB(dbUri, retryCount = 0) {
|
|
return new Promise((resolve, reject) => {
|
|
console.log('Connection Attempt : ', retryCount)
|
|
mongoose.set('strictQuery', true)
|
|
mongoose
|
|
.connect(dbUri)
|
|
.then((msg) => {
|
|
console.log('Connected to Mongo DB !')
|
|
resolve()
|
|
})
|
|
.catch((err) => {
|
|
console.log('Failed to connect dns mongo: ', err)
|
|
if (retryCount < 6) {
|
|
retryCount++
|
|
connectDB(dbUri, retryCount)
|
|
}
|
|
})
|
|
})
|
|
}
|
|
|
|
const processQueue = () => {
|
|
/* eslint-disable no-async-promise-executor */
|
|
return new Promise(async (resolve, reject) => {
|
|
try {
|
|
const response = await sqs.fetchMessageFromSQS(sqsQueueUrl)
|
|
|
|
if (
|
|
typeof response.Messages !== 'undefined' &&
|
|
response.Messages.length > 0
|
|
) {
|
|
throttleMessageFetching = false
|
|
const job = JSON.parse(response.Messages[0].Body)
|
|
const receiptHandle = response.Messages[0].ReceiptHandle
|
|
try {
|
|
await sqs.deleteMessageFromSQS(sqsQueueUrl, receiptHandle)
|
|
|
|
const {
|
|
userAudioProfileId,
|
|
text,
|
|
firstName,
|
|
salutationId,
|
|
recordingId,
|
|
baseUrlForPotionAi,
|
|
env,
|
|
} = job
|
|
|
|
const DB_URI =
|
|
env === 'production'
|
|
? mongoUriProd
|
|
: env === 'staging'
|
|
? mongoUriStaging
|
|
: mongoUriDev
|
|
|
|
console.log('DB_URI ', DB_URI)
|
|
await connectDB(DB_URI)
|
|
|
|
// read the path for the training model for the this users audio profile
|
|
|
|
const userAudioProfile = await userAudioProfileService.find({
|
|
_id: userAudioProfileId,
|
|
status: 'completed',
|
|
})
|
|
if (userAudioProfile) {
|
|
const { training_model_path, userId } = userAudioProfile[0]
|
|
const {
|
|
voice_model_light_path,
|
|
voice_model_config_light_path,
|
|
voice_model_speakers_file_path, // name for speakers embeddings file path
|
|
} = training_model_path
|
|
|
|
const outputPath = `/tmp/${uuid()}/`
|
|
if (!fs.existsSync(outputPath)) {
|
|
fs.mkdirSync(outputPath, { recursive: true })
|
|
}
|
|
|
|
const AI_COMMAND = `python3 ../voice-cloning/synthesize_speech.py --voice_model_path ${voice_model_light_path} --voice_model_config_path ${voice_model_config_light_path} --speaker_embeddings_path ${voice_model_speakers_file_path} --txt "${text}" --output_path ${outputPath}`
|
|
console.log('AI_COMMAND ', AI_COMMAND)
|
|
|
|
const SYNTHESIZE_AI_LABEL = `Time consumed by AI` + Math.random()
|
|
console.time(SYNTHESIZE_AI_LABEL)
|
|
const aiResponse = await execShellCommand(AI_COMMAND)
|
|
console.timeEnd(SYNTHESIZE_AI_LABEL)
|
|
|
|
let generatedFileName = ''
|
|
fs.readdirSync(`${outputPath}`).forEach((file) => {
|
|
if (file.includes('sr48000.wav')) generatedFileName = file
|
|
})
|
|
|
|
// upload the file to s3
|
|
const uploadParams = {
|
|
filePath: `${outputPath}${generatedFileName}`,
|
|
bucket: `recordings-${env}`,
|
|
fileName: `${uuid()}_salutation_${firstName.replace(
|
|
'-',
|
|
'_'
|
|
)}.wav`,
|
|
contentType: 'audio/x-wav',
|
|
fileType: 'wav',
|
|
}
|
|
console.time('Time to Upload video on S3')
|
|
const greetingUploadResponse = await s3.upload(uploadParams)
|
|
console.timeEnd('Time to Upload video on S3')
|
|
|
|
// Create new entry with the s3 path to salutation collection for the user and its profile id
|
|
// upsert the salutation
|
|
await salutationService.updateOrCreate(
|
|
{
|
|
firstName: firstName,
|
|
salutationVideo: greetingUploadResponse,
|
|
userAudioProfileId,
|
|
},
|
|
userId
|
|
)
|
|
// update the dynamic recordings for the current dynamic video with salutation url
|
|
const salutationToUpdate = await recordingSalutationModel.findOne({
|
|
_id: salutationId,
|
|
deleted: false,
|
|
})
|
|
|
|
const recordingToUpdate = await recordingModel.findOne({
|
|
_id: recordingId,
|
|
deleted: false,
|
|
})
|
|
|
|
if (
|
|
salutationToUpdate &&
|
|
salutationToUpdate.deleted === false &&
|
|
recordingToUpdate
|
|
) {
|
|
const jobsToInsert = []
|
|
|
|
await recordingSalutationModel.findOneAndUpdate(
|
|
{
|
|
_id: salutationId,
|
|
},
|
|
{
|
|
$set: {
|
|
salutationVideo: greetingUploadResponse,
|
|
},
|
|
}
|
|
)
|
|
|
|
const jobData = {
|
|
originalGreeting: recordingToUpdate.masterSalutationVideoUrl,
|
|
originalVideo:
|
|
recordingToUpdate.originalVideoUrl ||
|
|
recordingToUpdate.urls[0].url,
|
|
cropTimestamp: recordingToUpdate.cropTimestamp,
|
|
greetingClips: [greetingUploadResponse],
|
|
greetingObjects: [
|
|
{
|
|
greetingId: salutationToUpdate._id,
|
|
firstName: firstName,
|
|
videoUrl: greetingUploadResponse,
|
|
},
|
|
],
|
|
requestOrigin: baseUrlForPotionAi,
|
|
environment: env,
|
|
recordingId: recordingToUpdate._id,
|
|
salutation: salutationToUpdate._id,
|
|
dynamicVideoType: recordingToUpdate.dynamicVideoType,
|
|
}
|
|
jobsToInsert.push({
|
|
firstName,
|
|
recordingId: recordingToUpdate._id,
|
|
userId: recordingToUpdate.userId,
|
|
salutationId: salutationToUpdate._id,
|
|
metadata: jobData,
|
|
})
|
|
|
|
// create the job for the ai to create processing
|
|
if (jobsToInsert.length) {
|
|
await jobService.insertMany(jobsToInsert)
|
|
}
|
|
}
|
|
|
|
fs.unlinkSync(`${outputPath}${generatedFileName}`)
|
|
console.log(`[deleted] ${outputPath}${generatedFileName}`)
|
|
} else {
|
|
Bugsnag.notify(
|
|
new Error(
|
|
`audio profile training model not found ` + JSON.stringify(job)
|
|
)
|
|
)
|
|
|
|
resolve() // to continue working on new jobs
|
|
}
|
|
} catch (error) {
|
|
console.error('Error while synthesizing audio', { error })
|
|
Bugsnag.notify(
|
|
new Error(`Unable to synthesize audio ` + JSON.stringify(job))
|
|
)
|
|
Bugsnag.notify(error)
|
|
resolve() // to continue working on new jobs
|
|
}
|
|
} else {
|
|
throttleMessageFetching = true
|
|
}
|
|
resolve()
|
|
} catch (error) {
|
|
console.error('Error while synthesizing audio', { error })
|
|
Bugsnag.notify(error)
|
|
resolve() // to continue working on new jobs
|
|
} finally {
|
|
mongoose.connection.close()
|
|
}
|
|
})
|
|
}
|
|
|
|
function sleep(ms) {
|
|
return new Promise((resolve) => {
|
|
setTimeout(resolve, ms)
|
|
})
|
|
}
|
|
const init = async () => {
|
|
Bugsnag.start({
|
|
appVersion: APP_ENV + version,
|
|
apiKey: process.env.BUGSNAG_BACKEND_KEY,
|
|
releaseStage: process.env.NODE_ENV,
|
|
})
|
|
try {
|
|
while (true) {
|
|
await processQueue()
|
|
if (throttleMessageFetching) await sleep(2000)
|
|
}
|
|
} catch (error) {
|
|
Bugsnag.notify(error)
|
|
}
|
|
}
|
|
init()
|