Merge pull request #3926 from omnivore-app/feature/delete-task-api

feature/delete task api
This commit is contained in:
Hongbo Wu 2024-05-09 15:15:33 +08:00 committed by GitHub
commit ab85805e24
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
4 changed files with 129 additions and 23 deletions

View file

@ -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<express.Request>(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
}

View file

@ -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<express.Request>(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)
}
})

View file

@ -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))
}

View file

@ -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',
},
})