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()