diff --git a/packages/api/src/routers/digest_router.ts b/packages/api/src/routers/digest_router.ts index d136e9638..c47ce93aa 100644 --- a/packages/api/src/routers/digest_router.ts +++ b/packages/api/src/routers/digest_router.ts @@ -6,12 +6,12 @@ import { CreateDigestJobSchedule, moveDigestToLibrary, } from '../jobs/ai/create_digest' -import { getDigest } from '../services/digest' +import { deleteDigest, getDigest } from '../services/digest' import { FeatureName, findGrantedFeatureByName } from '../services/features' import { analytics } from '../utils/analytics' import { getClaimsByToken, getTokenByRequest } from '../utils/auth' import { corsConfig } from '../utils/corsConfig' -import { enqueueCreateDigest } from '../utils/createTask' +import { enqueueCreateDigest, removeDigestJobs } from '../utils/createTask' import { logger } from '../utils/logger' interface Feedback { @@ -93,6 +93,9 @@ export function digestRouter() { return res.status(202).send(digest) } + // remove existing digest jobs + await removeDigestJobs(userId) + // enqueue job and return job id const result = await enqueueCreateDigest( { @@ -315,5 +318,51 @@ export function digestRouter() { } ) + // v1 version of delete digest api + router.delete('/v1', cors(corsConfig), async (req, res) => { + const token = getTokenByRequest(req) + // get claims from token + const claims = await getClaimsByToken(token) + if (!claims) { + logger.error('Token not found') + return res.status(401).send({ + error: 'UNAUTHORIZED', + }) + } + + // get user by uid from claims + const userId = claims.uid + + try { + const feature = await findGrantedFeatureByName( + FeatureName.AIDigest, + userId + ) + if (!feature) { + logger.info(`${FeatureName.AIDigest} not granted: ${userId}`) + return res.status(403).send({ + error: 'FORBIDDEN', + }) + } + + // cancel and remove the digest job + await removeDigestJobs(userId) + logger.info(`Digest job removed: ${userId}`) + + // delete digest + await deleteDigest(userId) + logger.info(`Digest deleted: ${userId}`) + + res.send({ + success: true, + }) + } catch (error) { + logger.error('Error while deleting digest', error) + return res.status(500).send({ + error: 'INTERNAL_SERVER_ERROR', + }) + } + }) + return router } diff --git a/packages/api/src/routers/task_router.ts b/packages/api/src/routers/task_router.ts index f6e50d026..cd88d185a 100644 --- a/packages/api/src/routers/task_router.ts +++ b/packages/api/src/routers/task_router.ts @@ -41,7 +41,45 @@ export function taskRouter() { res.send(result) } catch (e) { logger.error('failed to get task', e) - res.status(500) + res.sendStatus(500) + } + }) + + router.delete('/:id', cors(corsConfig), async (req, res) => { + const token = getTokenByRequest(req) + const claims = await getClaimsByToken(token) + if (!claims) { + return res.sendStatus(401) + } + + try { + const job = await getJob(req.params.id) + if (!job || !job.id) { + logger.info('Task not found') + return res.sendStatus(404) + } + + const jobState = await job.getState() + if (jobState === 'active') { + logger.error('Task is active') + // cannot delete active task + return res.status(400).send('Task is active') + } + + // remove job + await job.remove() + + if (['completed', 'failed'].includes(jobState)) { + logger.info('Task removed') + return res.status(200).send('Task removed') + } + + // job is waiting or delayed + logger.info('Task cancelled') + res.status(200).send('Task cancelled') + } catch (e) { + logger.error('failed to delete task', e) + res.sendStatus(500) } }) diff --git a/packages/api/src/services/digest.ts b/packages/api/src/services/digest.ts index a5bc2537a..f80cbc1bb 100644 --- a/packages/api/src/services/digest.ts +++ b/packages/api/src/services/digest.ts @@ -51,3 +51,7 @@ export const writeDigest = async (userId: string, digest: Digest) => { throw new Error(msg) } } + +export const deleteDigest = async (userId: string) => { + await redisDataSource.redisClient?.del(digestKey(userId)) +} diff --git a/packages/api/src/utils/createTask.ts b/packages/api/src/utils/createTask.ts index d74f2ae10..47d1a7386 100644 --- a/packages/api/src/utils/createTask.ts +++ b/packages/api/src/utils/createTask.ts @@ -855,6 +855,40 @@ export const enqueueSendEmail = async (jobData: SendEmailJobData) => { }) } +export const scheduledDigestJobOptions = ( + schedule: CreateDigestJobSchedule +) => ({ + pattern: getCronPattern(schedule), + tz: 'UTC', +}) + +export const removeDigestJobs = async (userId: string) => { + const queue = await getBackendQueue() + if (!queue) { + throw new Error('No queue found') + } + + const jobId = `${CREATE_DIGEST_JOB}_${userId}` + + // remove existing one-time job if any + const job = await queue.getJob(jobId) + if (job) { + await job.remove() + logger.info('existing job removed', { jobId }) + } + + // remove existing repeated job if any + await Promise.all( + Object.keys(CRON_PATTERNS).map((key) => + queue.removeRepeatable( + CREATE_DIGEST_JOB, + scheduledDigestJobOptions(key as CreateDigestJobSchedule), + jobId + ) + ) + ) +} + export const enqueueCreateDigest = async ( data: CreateDigestData, schedule?: CreateDigestJobSchedule @@ -892,24 +926,6 @@ export const enqueueCreateDigest = async ( await writeDigest(data.userId, digest) if (schedule) { - // remove existing repeated job if any - await Promise.all( - Object.keys(CRON_PATTERNS).map(async (key) => { - const isDeleted = await queue.removeRepeatable( - CREATE_DIGEST_JOB, - { - pattern: CRON_PATTERNS[key as keyof typeof CRON_PATTERNS], - tz: 'UTC', - }, - jobId - ) - - if (isDeleted) { - logger.info('existing repeated job removed', { jobId, schedule: key }) - } - }) - ) - // schedule repeated job // delete the digest id to avoid duplication delete data.id @@ -918,9 +934,8 @@ export const enqueueCreateDigest = async ( attempts: 1, priority: getJobPriority(CREATE_DIGEST_JOB), repeat: { - pattern: getCronPattern(schedule), + ...scheduledDigestJobOptions(schedule), // cron parser options (tz, etc.) jobId, - tz: 'UTC', }, })