diff --git a/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/test/worker-tenant.test.js b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/test/worker-tenant.test.js new file mode 100644 index 0000000..025389e --- /dev/null +++ b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/test/worker-tenant.test.js @@ -0,0 +1,71 @@ +const test = require('node:test') +const assert = require('node:assert/strict') + +const { tenantFilter } = require('../worker_tenant') +const cloningService = require('../voice-cloning-job-handler/voice_cloning/voice_cloning_service') +const cloningProfileService = require('../voice-cloning-job-handler/user_audio_profile/user_audio_profile_service') +const synthProfileService = require('../voice-synthsizer-job-handler/user_audio_profile/user_audio_profile_service') +const salutationService = require('../voice-synthsizer-job-handler/salutation/salutation_service') +const jobService = require('../voice-synthsizer-job-handler/job/job_service') + +const userId = '507f1f77bcf86cd799439011' +const otherUserId = '507f191e810c19729de860ea' +const recordId = '507f1f77bcf86cd799439012' + +function fakeModel() { + const calls = [] + return { + calls, + findOne: async (filter) => { calls.push(['findOne', filter]); return null }, + find: async (filter) => { calls.push(['find', filter]); return [] }, + findOneAndUpdate: async (filter, update) => { + calls.push(['findOneAndUpdate', filter, update]) + return null + }, + updateMany: async (filter, update) => { + calls.push(['updateMany', filter, update]) + return { modifiedCount: 0 } + }, + } +} + +test('tenant filter requires a valid user ID and keeps it in the query', () => { + assert.throws(() => tenantFilter({ _id: recordId }), /userId/) + assert.throws(() => tenantFilter({ _id: recordId, userId: { $ne: otherUserId } }), /userId/) + assert.equal(String(tenantFilter({ _id: recordId, userId }).userId), userId) +}) + +for (const [name, factory] of [ + ['cloning record', cloningService], + ['cloning profile', cloningProfileService], + ['synthesis profile', synthProfileService], + ['synthesis job', jobService], +]) { + test(`${name} reads and writes require matching tenant`, async () => { + const model = fakeModel() + const service = factory(model) + await service.read({ _id: recordId, userId }) + await service.find({ userId }) + await service.update({ _id: recordId, userId, status: 'processing' }) + await service.remove({ _id: recordId, userId }) + await service.removeMany({ userId }) + for (const [, filter] of model.calls) { + assert.equal(String(filter.userId), userId) + } + const update = model.calls.find(([method]) => method === 'findOneAndUpdate') + assert.equal(update[1]._id, recordId) + assert.deepEqual(update[2], { $set: { status: 'processing' } }) + assert.equal(update[1].deleted, false) + await assert.rejects(service.update({ _id: recordId, status: 'processing' }), /userId/) + }) +} + +test('salutation updates cannot change ownership', async () => { + const model = fakeModel() + const service = salutationService(model) + await service.update({ _id: recordId, userId: otherUserId, salutationVideo: 'url' }, userId) + const [, filter, update] = model.calls[0] + assert.equal(String(filter.userId), userId) + assert.deepEqual(update, { $set: { salutationVideo: 'url' } }) + await assert.rejects(service.read({ _id: recordId }), /userId/) +}) diff --git a/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-cloning-job-handler/index.js b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-cloning-job-handler/index.js index 0995e23..84f050a 100644 --- a/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-cloning-job-handler/index.js +++ b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-cloning-job-handler/index.js @@ -10,6 +10,7 @@ const sqs = require('../app/services/sqs') const s3 = require('../app/services/s3') const voiceCloningService = require('./voice_cloning') const userAudioProfileService = require('./user_audio_profile') +const { requireUserId } = require('../worker_tenant') AWS.config.update({ region: 'us-west-2' }) const sqsQueueUrl = process.env.SQS_URL @@ -101,7 +102,7 @@ const processQueue = () => { const receiptHandle = response.Messages[0].ReceiptHandle console.log('job===', job) - const { metadata, input, _id, userAudioProfileId } = job._doc + const { metadata, input, _id, userAudioProfileId, userId } = job._doc console.log('userAudioProfileId', userAudioProfileId) console.log('_id', _id) const { env } = job @@ -126,9 +127,25 @@ const processQueue = () => { ? cloudFrontUrlStaging : cloudFrontUrlDev + let authorized = false try { await sqs.deleteMessageFromSQS(sqsQueueUrl, receiptHandle) + requireUserId(userId) + const cloningRecord = await voiceCloningService.read({ + _id, + userId, + userAudioProfileId, + }) + const audioProfile = await userAudioProfileService.read({ + _id: userAudioProfileId, + userId, + }) + if (!cloningRecord || !audioProfile) { + throw new Error('Voice cloning job records do not belong to the job user') + } + authorized = true + const { directoryName } = metadata console.log('directoryName', directoryName) const logPath = `/mnt/efs/potion-voice/${env}/${directoryName}` @@ -136,11 +153,19 @@ const processQueue = () => { fs.mkdirSync(logPath, { recursive: true }) } // update the db model to processing - await voiceCloningService.update({ _id, status: 'processing' }) - await userAudioProfileService.update({ - _id: userAudioProfileId, + const processingClone = await voiceCloningService.update({ + _id, + userId, status: 'processing', }) + const processingProfile = await userAudioProfileService.update({ + _id: userAudioProfileId, + userId, + status: 'processing', + }) + if (!processingClone || !processingProfile) { + throw new Error('Voice cloning job records are no longer available to the job user') + } // create directory for userid-useraudioprofileid if not exist const rootPath = `/tmp/${directoryName}` @@ -240,7 +265,7 @@ const processQueue = () => { console.timeEnd(VOICE_MINIMIZE_LABEL) // Add the code to update location of generated model and status into DB - await voiceCloningService.update({ _id, status: 'completed' }) + await voiceCloningService.update({ _id, userId, status: 'completed' }) const training_model_path = { voice_model_path: `${resultsPath}/${generatedDirectoryName}/checkpoint_365200.pth`, @@ -252,6 +277,7 @@ const processQueue = () => { await userAudioProfileService.update({ _id: userAudioProfileId, + userId, status: 'completed', training_model_path, }) @@ -273,6 +299,7 @@ const processQueue = () => { // add S3 path to user audio profile model await userAudioProfileService.update({ _id: userAudioProfileId, + userId, training_model_s3_path, }) } catch (error) { @@ -285,11 +312,14 @@ const processQueue = () => { Bugsnag.notify(error) // update the db to set status as error - await voiceCloningService.update({ _id, status: 'error' }) - await userAudioProfileService.update({ - _id: userAudioProfileId, - status: 'error', - }) + if (authorized) { + await voiceCloningService.update({ _id, userId, status: 'error' }) + await userAudioProfileService.update({ + _id: userAudioProfileId, + userId, + status: 'error', + }) + } resolve() // to continue working on new jobs } diff --git a/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-cloning-job-handler/user_audio_profile/user_audio_profile_service.js b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-cloning-job-handler/user_audio_profile/user_audio_profile_service.js index e76643d..244a848 100644 --- a/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-cloning-job-handler/user_audio_profile/user_audio_profile_service.js +++ b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-cloning-job-handler/user_audio_profile/user_audio_profile_service.js @@ -1,4 +1,5 @@ const StringifyUtils = require('../../app/services/utils/logService') +const { tenantFilter } = require('../../worker_tenant') const create = (UserAudioProfileModel) => async (data) => { try { @@ -32,7 +33,7 @@ const insertMany = (UserAudioProfileModel) => async (data) => { const read = (UserAudioProfileModel) => async (filter) => { try { const foundModel = await UserAudioProfileModel.findOne({ - ...filter, + ...tenantFilter(filter), deleted: false, }) return foundModel @@ -49,7 +50,7 @@ const read = (UserAudioProfileModel) => async (filter) => { const find = (UserAudioProfileModel) => async (filter) => { try { const foundModels = await UserAudioProfileModel.find({ - ...filter, + ...tenantFilter(filter), deleted: false, }) return foundModels @@ -65,9 +66,10 @@ const find = (UserAudioProfileModel) => async (filter) => { const update = (UserAudioProfileModel) => async (data) => { try { + const { _id, userId, ...changes } = data const updatedModel = await UserAudioProfileModel.findOneAndUpdate( - { _id: data._id }, - data, + { ...tenantFilter({ _id, userId }), deleted: false }, + { $set: changes }, { new: true, } @@ -86,7 +88,7 @@ const update = (UserAudioProfileModel) => async (data) => { const remove = (UserAudioProfileModel) => async (filter) => { try { const updatedModel = await UserAudioProfileModel.findOneAndUpdate( - { ...filter }, + { ...tenantFilter(filter) }, { $set: { deleted: true, @@ -108,7 +110,7 @@ const remove = (UserAudioProfileModel) => async (filter) => { const removeMany = (UserAudioProfileModel) => async (filter) => { try { const updatedModel = await UserAudioProfileModel.updateMany( - { ...filter }, + { ...tenantFilter(filter) }, { $set: { deleted: true, diff --git a/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-cloning-job-handler/voice_cloning/voice_cloning_service.js b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-cloning-job-handler/voice_cloning/voice_cloning_service.js index 329638e..8e7bd55 100644 --- a/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-cloning-job-handler/voice_cloning/voice_cloning_service.js +++ b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-cloning-job-handler/voice_cloning/voice_cloning_service.js @@ -1,4 +1,5 @@ const StringifyUtils = require('../../app/services/utils/logService') +const { tenantFilter } = require('../../worker_tenant') const create = (VoiceCloningModel) => async (data) => { try { @@ -32,7 +33,7 @@ const insertMany = (VoiceCloningModel) => async (data) => { const read = (VoiceCloningModel) => async (filter) => { try { const foundModel = await VoiceCloningModel.findOne({ - ...filter, + ...tenantFilter(filter), deleted: false, }) return foundModel @@ -49,7 +50,7 @@ const read = (VoiceCloningModel) => async (filter) => { const find = (VoiceCloningModel) => async (filter) => { try { const foundModels = await VoiceCloningModel.find({ - ...filter, + ...tenantFilter(filter), deleted: false, }) return foundModels @@ -65,9 +66,10 @@ const find = (VoiceCloningModel) => async (filter) => { const update = (VoiceCloningModel) => async (data) => { try { + const { _id, userId, ...changes } = data const updatedModel = await VoiceCloningModel.findOneAndUpdate( - { _id: data._id }, - data, + { ...tenantFilter({ _id, userId }), deleted: false }, + { $set: changes }, { new: true, } @@ -87,7 +89,7 @@ const update = (VoiceCloningModel) => async (data) => { const remove = (VoiceCloningModel) => async (filter) => { try { const updatedModel = await VoiceCloningModel.findOneAndUpdate( - { ...filter }, + { ...tenantFilter(filter) }, { $set: { deleted: true, @@ -109,7 +111,7 @@ const remove = (VoiceCloningModel) => async (filter) => { const removeMany = (VoiceCloningModel) => async (filter) => { try { const updatedModel = await VoiceCloningModel.updateMany( - { ...filter }, + { ...tenantFilter(filter) }, { $set: { deleted: true, diff --git a/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-synthsizer-job-handler/index.js b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-synthsizer-job-handler/index.js index 862300f..7cb19f0 100644 --- a/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-synthsizer-job-handler/index.js +++ b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-synthsizer-job-handler/index.js @@ -11,6 +11,7 @@ const recordingModel = require('./recording') const recordingSalutationModel = require('./recording_salutation') const jobService = require('./job') const salutationService = require('./salutation') +const { requireUserId } = require('../worker_tenant') let throttleMessageFetching = true AWS.config.update({ region: 'us-west-2' }) const sqsQueueUrl = process.env.SQS_URL @@ -72,6 +73,7 @@ const processQueue = () => { await sqs.deleteMessageFromSQS(sqsQueueUrl, receiptHandle) const { + userId, userAudioProfileId, text, firstName, @@ -91,14 +93,17 @@ const processQueue = () => { console.log('DB_URI ', DB_URI) await connectDB(DB_URI) + requireUserId(userId) + // read the path for the training model for the this users audio profile - const userAudioProfile = await userAudioProfileService.find({ + const userAudioProfile = await userAudioProfileService.read({ _id: userAudioProfileId, + userId, status: 'completed', }) if (userAudioProfile) { - const { training_model_path, userId } = userAudioProfile[0] + const { training_model_path } = userAudioProfile const { voice_model_light_path, voice_model_config_light_path, @@ -151,11 +156,13 @@ const processQueue = () => { // update the dynamic recordings for the current dynamic video with salutation url const salutationToUpdate = await recordingSalutationModel.findOne({ _id: salutationId, + userId, deleted: false, }) const recordingToUpdate = await recordingModel.findOne({ _id: recordingId, + userId, deleted: false, }) @@ -169,6 +176,8 @@ const processQueue = () => { await recordingSalutationModel.findOneAndUpdate( { _id: salutationId, + userId, + deleted: false, }, { $set: { diff --git a/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-synthsizer-job-handler/job/job_service.js b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-synthsizer-job-handler/job/job_service.js index 44189af..ff53150 100644 --- a/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-synthsizer-job-handler/job/job_service.js +++ b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-synthsizer-job-handler/job/job_service.js @@ -1,4 +1,5 @@ const StringifyUtils = require('../../app/services/utils/logService') +const { tenantFilter } = require('../../worker_tenant') const create = (Job) => async (jobData) => { try { @@ -32,7 +33,7 @@ const insertMany = (Job) => async (jobData) => { const read = (Job) => async (filter) => { try { const foundJob = await Job.findOne({ - ...filter, + ...tenantFilter(filter), deleted: false, }) return foundJob @@ -49,7 +50,7 @@ const read = (Job) => async (filter) => { const find = (Job) => async (filter) => { try { const foundJobs = await Job.find({ - ...filter, + ...tenantFilter(filter), deleted: false, }) return foundJobs @@ -65,9 +66,12 @@ const find = (Job) => async (filter) => { const update = (Job) => async (job) => { try { - const updatedJob = await Job.findOneAndUpdate({ _id: job._id }, job, { - new: true, - }) + const { _id, userId, ...changes } = job + const updatedJob = await Job.findOneAndUpdate( + { ...tenantFilter({ _id, userId }), deleted: false }, + { $set: changes }, + { new: true } + ) return updatedJob } catch (error) { const details = { job } @@ -82,7 +86,7 @@ const update = (Job) => async (job) => { const remove = (Job) => async (filter) => { try { const updatedJob = await Job.findOneAndUpdate( - { ...filter }, + { ...tenantFilter(filter) }, { $set: { deleted: true, @@ -104,7 +108,7 @@ const remove = (Job) => async (filter) => { const removeMany = (Job) => async (filter) => { try { const updatedJob = await Job.updateMany( - { ...filter }, + { ...tenantFilter(filter) }, { $set: { deleted: true, diff --git a/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-synthsizer-job-handler/salutation/salutation_service.js b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-synthsizer-job-handler/salutation/salutation_service.js index 2d84f78..ad58f0d 100644 --- a/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-synthsizer-job-handler/salutation/salutation_service.js +++ b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-synthsizer-job-handler/salutation/salutation_service.js @@ -1,4 +1,7 @@ +const { tenantFilter, requireUserId } = require('../../worker_tenant') + const create = (Salutation) => async (salutationData, userId) => { + requireUserId(userId) const newSalutation = new Salutation({ ...salutationData, userId }) const savedSalutation = await newSalutation.save() return savedSalutation @@ -6,7 +9,7 @@ const create = (Salutation) => async (salutationData, userId) => { const read = (Salutation) => async (filter) => { const foundSalutation = await Salutation.findOne({ - ...filter, + ...tenantFilter(filter), deleted: false, }) return foundSalutation @@ -14,29 +17,33 @@ const read = (Salutation) => async (filter) => { const find = (Salutation) => async (filter) => { const foundSalutations = await Salutation.find({ - ...filter, + ...tenantFilter(filter), deleted: false, }) return foundSalutations } const update = (Salutation) => async (salutation, userId) => { + const { _id, userId: ignoredUserId, ...changes } = salutation.toObject + ? salutation.toObject() + : salutation const updatedSalutation = await Salutation.findOneAndUpdate( - { _id: salutation._id, userId }, - salutation, + { ...tenantFilter({ _id, userId }), deleted: false }, + { $set: changes }, { new: true } ) return updatedSalutation } -const modify = (Salutation) => async (salutation) => { - const findQuery = salutation._id - ? { _id: salutation._id } +const modify = (Salutation) => async (salutation, userId) => { + const { _id, userId: ignoredUserId, ...changes } = salutation + const findQuery = _id + ? { _id } : { transcriptId: salutation.transcriptId } const updatedSalutation = await Salutation.findOneAndUpdate( - findQuery, - salutation, + { ...tenantFilter({ ...findQuery, userId }), deleted: false }, + { $set: changes }, { new: true } ) return updatedSalutation @@ -63,7 +70,7 @@ const updateOrCreate = const remove = (Salutation) => async (filter, userId) => { const updatedSalutation = await Salutation.findOneAndUpdate( - { ...filter, userId }, + { ...tenantFilter({ ...filter, userId }), deleted: false }, { $set: { deleted: true, diff --git a/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-synthsizer-job-handler/user_audio_profile/user_audio_profile_service.js b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-synthsizer-job-handler/user_audio_profile/user_audio_profile_service.js index 18416d7..22acbf8 100644 --- a/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-synthsizer-job-handler/user_audio_profile/user_audio_profile_service.js +++ b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/voice-synthsizer-job-handler/user_audio_profile/user_audio_profile_service.js @@ -1,4 +1,5 @@ const StringifyUtils = require('../../app/services/utils/logService') +const { tenantFilter } = require('../../worker_tenant') const create = (UserAudioProfileModel) => async (data) => { try { @@ -32,7 +33,7 @@ const insertMany = (UserAudioProfileModel) => async (data) => { const read = (UserAudioProfileModel) => async (filter) => { try { const foundModel = await UserAudioProfileModel.findOne({ - ...filter, + ...tenantFilter(filter), deleted: false }) return foundModel @@ -49,7 +50,7 @@ const read = (UserAudioProfileModel) => async (filter) => { const find = (UserAudioProfileModel) => async (filter) => { try { const foundModels = await UserAudioProfileModel.find({ - ...filter, + ...tenantFilter(filter), deleted: false }) return foundModels @@ -66,9 +67,10 @@ const find = (UserAudioProfileModel) => async (filter) => { const update = (UserAudioProfileModel) => async (data) => { console.log('ua data', data) try { + const { _id, userId, ...changes } = data const updatedModel = await UserAudioProfileModel.findOneAndUpdate( - { _id: data._id }, - data, + { ...tenantFilter({ _id, userId }), deleted: false }, + { $set: changes }, { new: true } @@ -88,7 +90,7 @@ const update = (UserAudioProfileModel) => async (data) => { const remove = (UserAudioProfileModel) => async (filter) => { try { const updatedModel = await UserAudioProfileModel.findOneAndUpdate( - { ...filter }, + { ...tenantFilter(filter) }, { $set: { deleted: true @@ -110,7 +112,7 @@ const remove = (UserAudioProfileModel) => async (filter) => { const removeMany = (UserAudioProfileModel) => async (filter) => { try { const updatedModel = await UserAudioProfileModel.updateMany( - { ...filter }, + { ...tenantFilter(filter) }, { $set: { deleted: true diff --git a/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/worker_tenant.js b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/worker_tenant.js new file mode 100644 index 0000000..f599574 --- /dev/null +++ b/worker-toolkit-potion-polyglot-v1.0.1/repos/potion-voice/worker_tenant.js @@ -0,0 +1,18 @@ +const mongoose = require('mongoose') + +const requireUserId = (userId) => { + if (typeof userId !== 'string' || !/^[a-f\d]{24}$/i.test(userId)) { + throw new Error('A valid job userId is required') + } + return userId +} + +const tenantFilter = (filter) => { + if (!filter || typeof filter !== 'object') { + throw new Error('A tenant-scoped filter is required') + } + const userId = requireUserId(filter.userId) + return { ...filter, userId: new mongoose.Types.ObjectId(userId) } +} + +module.exports = { requireUserId, tenantFilter }