Files
2026-10-05 16:14:53 -04:00

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