From 916fe7ace383b4da169e8bf5d75f7aed6961e70d Mon Sep 17 00:00:00 2001 From: Eric Bell Date: Thu, 8 Oct 2026 14:55:25 -0400 Subject: [PATCH] changes to implement instructions --- .../potion-voice/test/worker-tenant.test.js | 132 ++++++++++++++++++ .../voice-cloning-job-handler/index.js | 45 ++++-- .../user_audio_profile_service.js | 22 +-- .../voice_cloning/voice_cloning_service.js | 22 +-- .../voice-synthsizer-job-handler/index.js | 39 ++++-- .../job/job_service.js | 22 ++- .../salutation/salutation_service.js | 29 ++-- .../user_audio_profile_service.js | 22 +-- .../repos/potion-voice/worker-tenant.js | 12 ++ 9 files changed, 279 insertions(+), 66 deletions(-) create mode 100644 worker-toolkit-potion-polyglot/repos/potion-voice/test/worker-tenant.test.js create mode 100644 worker-toolkit-potion-polyglot/repos/potion-voice/worker-tenant.js diff --git a/worker-toolkit-potion-polyglot/repos/potion-voice/test/worker-tenant.test.js b/worker-toolkit-potion-polyglot/repos/potion-voice/test/worker-tenant.test.js new file mode 100644 index 0000000..9f742ae --- /dev/null +++ b/worker-toolkit-potion-polyglot/repos/potion-voice/test/worker-tenant.test.js @@ -0,0 +1,132 @@ +const assert = require('node:assert/strict') +const { test } = require('node:test') + +const userId = '507f1f77bcf86cd799439011' +const otherUserId = '507f191e810c19729de860ea' +const { requireDocumentId, requireUserId } = require('../worker-tenant') + +test('queue IDs reject missing values and MongoDB query operators', () => { + assert.equal(requireUserId(userId), userId) + assert.equal(requireDocumentId(otherUserId), otherUserId) + assert.throws(() => requireUserId({ $ne: otherUserId }), /userId/) + assert.throws(() => requireDocumentId(undefined), /document ID/) +}) + +const factories = [ + require('../voice-cloning-job-handler/voice_cloning/voice_cloning_service'), + require('../voice-cloning-job-handler/user_audio_profile/user_audio_profile_service'), + require('../voice-synthsizer-job-handler/user_audio_profile/user_audio_profile_service'), + require('../voice-synthsizer-job-handler/job/job_service'), +] + +for (const [index, factory] of factories.entries()) { + test(`worker service ${index} scopes reads and writes`, async () => { + const calls = [] + class Model { + static async findOne(filter) { calls.push(['findOne', filter]) } + static async find(filter) { calls.push(['find', filter]); return [] } + static async findOneAndUpdate(filter, update) { + calls.push(['findOneAndUpdate', filter, update]) + } + static async updateMany(filter, update) { + calls.push(['updateMany', filter, update]) + } + static async insertMany(data) { calls.push(['insertMany', data]) } + } + const service = factory(Model) + + await service.read({ _id: 'record', userId }) + await service.find({ userId }) + await service.update({ _id: 'record', userId, status: 'completed' }) + await service.remove({ _id: 'record', userId }) + await service.removeMany({ userId }) + await service.insertMany([{ userId, status: 'created' }]) + + for (const call of calls.slice(0, 5)) { + assert.equal(call[1].userId, userId) + } + assert.equal(calls[2][1].deleted, false) + assert.deepEqual(calls[2][2], { $set: { status: 'completed' } }) + assert.equal(calls[5][1][0].userId, userId) + + const originalLog = console.log + const before = calls.length + console.log = () => {} + try { + await assert.rejects(service.read({ _id: 'record' }), /userId/) + await assert.rejects( + service.update({ _id: 'record', userId: { $ne: otherUserId } }), + /userId/ + ) + await assert.rejects(service.remove({ _id: 'record' }), /userId/) + await assert.rejects( + service.insertMany([{ userId }, { status: 'created' }]), + /userId/ + ) + } finally { + console.log = originalLog + } + assert.equal(calls.length, before) + }) +} + +test('salutation upsert stays within the job tenant', async () => { + const calls = [] + class Model { + constructor(data) { calls.push(['create', data]) } + async save() { return this } + static async findOne(filter) { calls.push(['findOne', filter]); return null } + } + const service = require('../voice-synthsizer-job-handler/salutation/salutation_service')(Model) + await service.updateOrCreate({ + firstName: 'Ada', + salutationVideo: 'url', + userAudioProfileId: 'profile', + }, userId) + assert.equal(calls[0][1].userId, userId) + assert.equal(calls[0][1].deleted, false) + assert.equal(calls[1][1].userId, userId) +}) + +test('salutation upsert updates an existing record within the job tenant', async () => { + const calls = [] + class Model { + static async findOne(filter) { + calls.push(['findOne', filter]) + return { _id: 'salutation', userId } + } + static async findOneAndUpdate(filter, update) { + calls.push(['findOneAndUpdate', filter, update]) + } + } + const service = require('../voice-synthsizer-job-handler/salutation/salutation_service')(Model) + await service.updateOrCreate({ + firstName: 'Ada', + salutationVideo: 'url', + userAudioProfileId: 'profile', + }, userId) + assert.equal(calls[0][1].userId, userId) + assert.deepEqual(calls[1][1], { _id: 'salutation', userId, deleted: false }) + assert.deepEqual(calls[1][2], { $set: { salutationVideo: 'url' } }) +}) + +test('salutation updates cannot change ownership', async () => { + const calls = [] + class Model { + static async findOneAndUpdate(filter, update) { + calls.push([filter, update]) + } + } + const service = require('../voice-synthsizer-job-handler/salutation/salutation_service')(Model) + await service.update({ + _id: 'salutation', + userId: otherUserId, + salutationVideo: 'url', + }, userId) + assert.deepEqual(calls[0], [ + { _id: 'salutation', userId, deleted: false }, + { $set: { salutationVideo: 'url' } }, + ]) + await assert.rejects(service.modify({ _id: 'salutation' }), /userId/) + assert.equal(calls.length, 1) +}) diff --git a/worker-toolkit-potion-polyglot/repos/potion-voice/voice-cloning-job-handler/index.js b/worker-toolkit-potion-polyglot/repos/potion-voice/voice-cloning-job-handler/index.js index 0995e23..d3552ca 100644 --- a/worker-toolkit-potion-polyglot/repos/potion-voice/voice-cloning-job-handler/index.js +++ b/worker-toolkit-potion-polyglot/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, requireDocumentId } = require('../worker-tenant') AWS.config.update({ region: 'us-west-2' }) const sqsQueueUrl = process.env.SQS_URL @@ -101,14 +102,15 @@ const processQueue = () => { const receiptHandle = response.Messages[0].ReceiptHandle console.log('job===', job) - const { metadata, input, _id, userAudioProfileId } = job._doc + const { _id, userAudioProfileId, userId } = job._doc + requireUserId(userId) + requireDocumentId(_id) + requireDocumentId(userAudioProfileId) console.log('userAudioProfileId', userAudioProfileId) console.log('_id', _id) const { env } = job console.log('env', env) - console.log('metadata------', metadata) - console.log('input', input) const DB_URI = env === 'production' ? mongoUriProd @@ -126,9 +128,26 @@ const processQueue = () => { ? cloudFrontUrlStaging : cloudFrontUrlDev + let authorized = false try { await sqs.deleteMessageFromSQS(sqsQueueUrl, receiptHandle) + const cloningJob = await voiceCloningService.read({ + _id, + userId, + userAudioProfileId, + }) + const audioProfile = await userAudioProfileService.read({ + _id: userAudioProfileId, + userId, + }) + if (!cloningJob || !audioProfile) { + throw new Error('Cloning job or audio profile does not belong to user') + } + authorized = true + + const { metadata, input } = cloningJob + const { directoryName } = metadata console.log('directoryName', directoryName) const logPath = `/mnt/efs/potion-voice/${env}/${directoryName}` @@ -136,9 +155,10 @@ const processQueue = () => { fs.mkdirSync(logPath, { recursive: true }) } // update the db model to processing - await voiceCloningService.update({ _id, status: 'processing' }) + await voiceCloningService.update({ _id, userId, status: 'processing' }) await userAudioProfileService.update({ _id: userAudioProfileId, + userId, status: 'processing', }) @@ -240,7 +260,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 +272,7 @@ const processQueue = () => { await userAudioProfileService.update({ _id: userAudioProfileId, + userId, status: 'completed', training_model_path, }) @@ -273,6 +294,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 +307,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/repos/potion-voice/voice-cloning-job-handler/user_audio_profile/user_audio_profile_service.js b/worker-toolkit-potion-polyglot/repos/potion-voice/voice-cloning-job-handler/user_audio_profile/user_audio_profile_service.js index e76643d..f35ecd6 100644 --- a/worker-toolkit-potion-polyglot/repos/potion-voice/voice-cloning-job-handler/user_audio_profile/user_audio_profile_service.js +++ b/worker-toolkit-potion-polyglot/repos/potion-voice/voice-cloning-job-handler/user_audio_profile/user_audio_profile_service.js @@ -1,8 +1,9 @@ +const { requireUserId } = require('../../worker-tenant') const StringifyUtils = require('../../app/services/utils/logService') const create = (UserAudioProfileModel) => async (data) => { try { - const newModel = new UserAudioProfileModel({ ...data }) + const newModel = new UserAudioProfileModel({ ...data, userId: requireUserId(data.userId) }) const savedModel = await newModel.save() return savedModel } catch (error) { @@ -17,7 +18,9 @@ const create = (UserAudioProfileModel) => async (data) => { const insertMany = (UserAudioProfileModel) => async (data) => { try { - const inserted = await UserAudioProfileModel.insertMany(data) + const inserted = await UserAudioProfileModel.insertMany( + data.map((item) => ({ ...item, userId: requireUserId(item.userId) })) + ) return inserted } catch (error) { const details = { data } @@ -33,6 +36,7 @@ const read = (UserAudioProfileModel) => async (filter) => { try { const foundModel = await UserAudioProfileModel.findOne({ ...filter, + userId: requireUserId(filter.userId), deleted: false, }) return foundModel @@ -50,6 +54,7 @@ const find = (UserAudioProfileModel) => async (filter) => { try { const foundModels = await UserAudioProfileModel.find({ ...filter, + userId: requireUserId(filter.userId), deleted: false, }) return foundModels @@ -65,12 +70,11 @@ 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, - { - new: true, - } + { _id, userId: requireUserId(userId), deleted: false }, + { $set: changes }, + { new: true } ) return updatedModel } catch (error) { @@ -86,7 +90,7 @@ const update = (UserAudioProfileModel) => async (data) => { const remove = (UserAudioProfileModel) => async (filter) => { try { const updatedModel = await UserAudioProfileModel.findOneAndUpdate( - { ...filter }, + { ...filter, userId: requireUserId(filter.userId) }, { $set: { deleted: true, @@ -108,7 +112,7 @@ const remove = (UserAudioProfileModel) => async (filter) => { const removeMany = (UserAudioProfileModel) => async (filter) => { try { const updatedModel = await UserAudioProfileModel.updateMany( - { ...filter }, + { ...filter, userId: requireUserId(filter.userId) }, { $set: { deleted: true, diff --git a/worker-toolkit-potion-polyglot/repos/potion-voice/voice-cloning-job-handler/voice_cloning/voice_cloning_service.js b/worker-toolkit-potion-polyglot/repos/potion-voice/voice-cloning-job-handler/voice_cloning/voice_cloning_service.js index 329638e..cef601b 100644 --- a/worker-toolkit-potion-polyglot/repos/potion-voice/voice-cloning-job-handler/voice_cloning/voice_cloning_service.js +++ b/worker-toolkit-potion-polyglot/repos/potion-voice/voice-cloning-job-handler/voice_cloning/voice_cloning_service.js @@ -1,8 +1,9 @@ +const { requireUserId } = require('../../worker-tenant') const StringifyUtils = require('../../app/services/utils/logService') const create = (VoiceCloningModel) => async (data) => { try { - const newModel = new VoiceCloningModel({ ...data }) + const newModel = new VoiceCloningModel({ ...data, userId: requireUserId(data.userId) }) const savedModel = await newModel.save() return savedModel } catch (error) { @@ -17,7 +18,9 @@ const create = (VoiceCloningModel) => async (data) => { const insertMany = (VoiceCloningModel) => async (data) => { try { - const inserted = await VoiceCloningModel.insertMany(data) + const inserted = await VoiceCloningModel.insertMany( + data.map((item) => ({ ...item, userId: requireUserId(item.userId) })) + ) return inserted } catch (error) { const details = { data } @@ -33,6 +36,7 @@ const read = (VoiceCloningModel) => async (filter) => { try { const foundModel = await VoiceCloningModel.findOne({ ...filter, + userId: requireUserId(filter.userId), deleted: false, }) return foundModel @@ -50,6 +54,7 @@ const find = (VoiceCloningModel) => async (filter) => { try { const foundModels = await VoiceCloningModel.find({ ...filter, + userId: requireUserId(filter.userId), deleted: false, }) return foundModels @@ -65,12 +70,11 @@ 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, - { - new: true, - } + { _id, userId: requireUserId(userId), deleted: false }, + { $set: changes }, + { new: true } ) return updatedModel @@ -87,7 +91,7 @@ const update = (VoiceCloningModel) => async (data) => { const remove = (VoiceCloningModel) => async (filter) => { try { const updatedModel = await VoiceCloningModel.findOneAndUpdate( - { ...filter }, + { ...filter, userId: requireUserId(filter.userId) }, { $set: { deleted: true, @@ -109,7 +113,7 @@ const remove = (VoiceCloningModel) => async (filter) => { const removeMany = (VoiceCloningModel) => async (filter) => { try { const updatedModel = await VoiceCloningModel.updateMany( - { ...filter }, + { ...filter, userId: requireUserId(filter.userId) }, { $set: { deleted: true, diff --git a/worker-toolkit-potion-polyglot/repos/potion-voice/voice-synthsizer-job-handler/index.js b/worker-toolkit-potion-polyglot/repos/potion-voice/voice-synthsizer-job-handler/index.js index 862300f..7269ea7 100644 --- a/worker-toolkit-potion-polyglot/repos/potion-voice/voice-synthsizer-job-handler/index.js +++ b/worker-toolkit-potion-polyglot/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, requireDocumentId } = require('../worker-tenant') let throttleMessageFetching = true AWS.config.update({ region: 'us-west-2' }) const sqsQueueUrl = process.env.SQS_URL @@ -79,8 +80,14 @@ const processQueue = () => { recordingId, baseUrlForPotionAi, env, + userId, } = job + requireUserId(userId) + requireDocumentId(userAudioProfileId) + requireDocumentId(salutationId) + requireDocumentId(recordingId) + const DB_URI = env === 'production' ? mongoUriProd @@ -95,10 +102,24 @@ const processQueue = () => { const userAudioProfile = await userAudioProfileService.find({ _id: userAudioProfileId, + userId, status: 'completed', }) - if (userAudioProfile) { - const { training_model_path, userId } = userAudioProfile[0] + if (userAudioProfile.length) { + const { training_model_path } = userAudioProfile[0] + const salutationToUpdate = await recordingSalutationModel.findOne({ + _id: salutationId, + userId, + deleted: false, + }) + const recordingToUpdate = await recordingModel.findOne({ + _id: recordingId, + userId, + deleted: false, + }) + if (!salutationToUpdate || !recordingToUpdate) { + throw new Error('Job recording or salutation does not belong to user') + } const { voice_model_light_path, voice_model_config_light_path, @@ -149,16 +170,6 @@ const processQueue = () => { 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 && @@ -169,6 +180,8 @@ const processQueue = () => { await recordingSalutationModel.findOneAndUpdate( { _id: salutationId, + userId, + deleted: false, }, { $set: { @@ -200,7 +213,7 @@ const processQueue = () => { jobsToInsert.push({ firstName, recordingId: recordingToUpdate._id, - userId: recordingToUpdate.userId, + userId, salutationId: salutationToUpdate._id, metadata: jobData, }) diff --git a/worker-toolkit-potion-polyglot/repos/potion-voice/voice-synthsizer-job-handler/job/job_service.js b/worker-toolkit-potion-polyglot/repos/potion-voice/voice-synthsizer-job-handler/job/job_service.js index 44189af..290b1c4 100644 --- a/worker-toolkit-potion-polyglot/repos/potion-voice/voice-synthsizer-job-handler/job/job_service.js +++ b/worker-toolkit-potion-polyglot/repos/potion-voice/voice-synthsizer-job-handler/job/job_service.js @@ -1,8 +1,9 @@ +const { requireUserId } = require('../../worker-tenant') const StringifyUtils = require('../../app/services/utils/logService') const create = (Job) => async (jobData) => { try { - const newJob = new Job({ ...jobData }) + const newJob = new Job({ ...jobData, userId: requireUserId(jobData.userId) }) const savedJob = await newJob.save() return savedJob } catch (error) { @@ -17,7 +18,9 @@ const create = (Job) => async (jobData) => { const insertMany = (Job) => async (jobData) => { try { - const inserted = await Job.insertMany(jobData) + const inserted = await Job.insertMany( + jobData.map((item) => ({ ...item, userId: requireUserId(item.userId) })) + ) return inserted } catch (error) { const details = { jobData } @@ -33,6 +36,7 @@ const read = (Job) => async (filter) => { try { const foundJob = await Job.findOne({ ...filter, + userId: requireUserId(filter.userId), deleted: false, }) return foundJob @@ -50,6 +54,7 @@ const find = (Job) => async (filter) => { try { const foundJobs = await Job.find({ ...filter, + userId: requireUserId(filter.userId), deleted: false, }) return foundJobs @@ -65,9 +70,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( + { _id, userId: requireUserId(userId), deleted: false }, + { $set: changes }, + { new: true } + ) return updatedJob } catch (error) { const details = { job } @@ -82,7 +90,7 @@ const update = (Job) => async (job) => { const remove = (Job) => async (filter) => { try { const updatedJob = await Job.findOneAndUpdate( - { ...filter }, + { ...filter, userId: requireUserId(filter.userId) }, { $set: { deleted: true, @@ -104,7 +112,7 @@ const remove = (Job) => async (filter) => { const removeMany = (Job) => async (filter) => { try { const updatedJob = await Job.updateMany( - { ...filter }, + { ...filter, userId: requireUserId(filter.userId) }, { $set: { deleted: true, diff --git a/worker-toolkit-potion-polyglot/repos/potion-voice/voice-synthsizer-job-handler/salutation/salutation_service.js b/worker-toolkit-potion-polyglot/repos/potion-voice/voice-synthsizer-job-handler/salutation/salutation_service.js index 2d84f78..29dfc65 100644 --- a/worker-toolkit-potion-polyglot/repos/potion-voice/voice-synthsizer-job-handler/salutation/salutation_service.js +++ b/worker-toolkit-potion-polyglot/repos/potion-voice/voice-synthsizer-job-handler/salutation/salutation_service.js @@ -1,5 +1,7 @@ +const { requireUserId } = require('../../worker-tenant') + const create = (Salutation) => async (salutationData, userId) => { - const newSalutation = new Salutation({ ...salutationData, userId }) + const newSalutation = new Salutation({ ...salutationData, userId: requireUserId(userId) }) const savedSalutation = await newSalutation.save() return savedSalutation } @@ -7,6 +9,7 @@ const create = (Salutation) => async (salutationData, userId) => { const read = (Salutation) => async (filter) => { const foundSalutation = await Salutation.findOne({ ...filter, + userId: requireUserId(filter.userId), deleted: false, }) return foundSalutation @@ -15,28 +18,34 @@ const read = (Salutation) => async (filter) => { const find = (Salutation) => async (filter) => { const foundSalutations = await Salutation.find({ ...filter, + userId: requireUserId(filter.userId), deleted: false, }) return foundSalutations } const update = (Salutation) => async (salutation, userId) => { + const { _id, ...changes } = salutation + delete changes.userId const updatedSalutation = await Salutation.findOneAndUpdate( - { _id: salutation._id, userId }, - salutation, + { _id, userId: requireUserId(userId), deleted: false }, + { $set: changes }, { new: true } ) return updatedSalutation } -const modify = (Salutation) => async (salutation) => { +const modify = (Salutation) => async (salutation, userId) => { const findQuery = salutation._id ? { _id: salutation._id } : { transcriptId: salutation.transcriptId } + const changes = { ...salutation } + delete changes._id + delete changes.userId const updatedSalutation = await Salutation.findOneAndUpdate( - findQuery, - salutation, + { ...findQuery, userId: requireUserId(userId), deleted: false }, + { $set: changes }, { new: true } ) return updatedSalutation @@ -51,8 +60,10 @@ const updateOrCreate = userId, }) if (salutation) { - salutation.salutationVideo = salutationVideo - return await update(Salutation)(salutation, userId) + return await update(Salutation)( + { _id: salutation._id, salutationVideo }, + userId + ) } else { return await create(Salutation)( { firstName, salutationVideo, userAudioProfileId }, @@ -63,7 +74,7 @@ const updateOrCreate = const remove = (Salutation) => async (filter, userId) => { const updatedSalutation = await Salutation.findOneAndUpdate( - { ...filter, userId }, + { ...filter, userId: requireUserId(userId) }, { $set: { deleted: true, diff --git a/worker-toolkit-potion-polyglot/repos/potion-voice/voice-synthsizer-job-handler/user_audio_profile/user_audio_profile_service.js b/worker-toolkit-potion-polyglot/repos/potion-voice/voice-synthsizer-job-handler/user_audio_profile/user_audio_profile_service.js index 18416d7..a432cfc 100644 --- a/worker-toolkit-potion-polyglot/repos/potion-voice/voice-synthsizer-job-handler/user_audio_profile/user_audio_profile_service.js +++ b/worker-toolkit-potion-polyglot/repos/potion-voice/voice-synthsizer-job-handler/user_audio_profile/user_audio_profile_service.js @@ -1,8 +1,9 @@ +const { requireUserId } = require('../../worker-tenant') const StringifyUtils = require('../../app/services/utils/logService') const create = (UserAudioProfileModel) => async (data) => { try { - const newModel = new UserAudioProfileModel({ ...data }) + const newModel = new UserAudioProfileModel({ ...data, userId: requireUserId(data.userId) }) const savedModel = await newModel.save() return savedModel } catch (error) { @@ -17,7 +18,9 @@ const create = (UserAudioProfileModel) => async (data) => { const insertMany = (UserAudioProfileModel) => async (data) => { try { - const inserted = await UserAudioProfileModel.insertMany(data) + const inserted = await UserAudioProfileModel.insertMany( + data.map((item) => ({ ...item, userId: requireUserId(item.userId) })) + ) return inserted } catch (error) { const details = { data } @@ -33,6 +36,7 @@ const read = (UserAudioProfileModel) => async (filter) => { try { const foundModel = await UserAudioProfileModel.findOne({ ...filter, + userId: requireUserId(filter.userId), deleted: false }) return foundModel @@ -50,6 +54,7 @@ const find = (UserAudioProfileModel) => async (filter) => { try { const foundModels = await UserAudioProfileModel.find({ ...filter, + userId: requireUserId(filter.userId), deleted: false }) return foundModels @@ -66,12 +71,11 @@ 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, - { - new: true - } + { _id, userId: requireUserId(userId), deleted: false }, + { $set: changes }, + { new: true } ) console.log('ua updatedModel', updatedModel) return updatedModel @@ -88,7 +92,7 @@ const update = (UserAudioProfileModel) => async (data) => { const remove = (UserAudioProfileModel) => async (filter) => { try { const updatedModel = await UserAudioProfileModel.findOneAndUpdate( - { ...filter }, + { ...filter, userId: requireUserId(filter.userId) }, { $set: { deleted: true @@ -110,7 +114,7 @@ const remove = (UserAudioProfileModel) => async (filter) => { const removeMany = (UserAudioProfileModel) => async (filter) => { try { const updatedModel = await UserAudioProfileModel.updateMany( - { ...filter }, + { ...filter, userId: requireUserId(filter.userId) }, { $set: { deleted: true diff --git a/worker-toolkit-potion-polyglot/repos/potion-voice/worker-tenant.js b/worker-toolkit-potion-polyglot/repos/potion-voice/worker-tenant.js new file mode 100644 index 0000000..2b8c40a --- /dev/null +++ b/worker-toolkit-potion-polyglot/repos/potion-voice/worker-tenant.js @@ -0,0 +1,12 @@ +// SQS payloads must never turn a tenant or document filter into a broad query. +const requireId = (id, name) => { + if (typeof id !== 'string' || !/^[0-9a-f]{24}$/i.test(id)) { + throw new Error(`A valid ${name} is required`) + } + return id +} + +const requireUserId = (userId) => requireId(userId, 'userId') +const requireDocumentId = (id) => requireId(id, 'document ID') + +module.exports = { requireUserId, requireDocumentId }