From fd4f7fee1cf11a5fad2b235260dc6b6f5587b40f Mon Sep 17 00:00:00 2001 From: Junyi Date: Sun, 28 Sep 2025 20:16:18 +0800 Subject: [PATCH] fix(plugin-workflow): fix schedule trigger on date field not run when start (#7524) * fix(plugin-workflow): fix schedule trigger on date field not run when start * chore(plugin-workflow): remove console logs * fix(plugin-workflow): fix default unit of on field to days * fix(plugin-workflow): fix load condition * fix(plugin-workflow): fix mod expression * fix(plugin-workflow): fix default unit --- .../triggers/schedule/mode-date-field.test.ts | 82 +++++++++++++++++++ .../DateFieldScheduleTrigger.ts | 76 +++++++++++------ 2 files changed, 131 insertions(+), 27 deletions(-) diff --git a/packages/plugins/@nocobase/plugin-workflow/src/server/__tests__/triggers/schedule/mode-date-field.test.ts b/packages/plugins/@nocobase/plugin-workflow/src/server/__tests__/triggers/schedule/mode-date-field.test.ts index c67475f6db7..63303b073bb 100644 --- a/packages/plugins/@nocobase/plugin-workflow/src/server/__tests__/triggers/schedule/mode-date-field.test.ts +++ b/packages/plugins/@nocobase/plugin-workflow/src/server/__tests__/triggers/schedule/mode-date-field.test.ts @@ -78,6 +78,7 @@ describe('workflow > triggers > schedule > date field mode', () => { startsOn: { field: 'createdAt', offset: 2, + unit: 1000, }, }, }); @@ -110,6 +111,7 @@ describe('workflow > triggers > schedule > date field mode', () => { startsOn: { field: 'createdAt', offset: -2, + unit: 1000, }, }, }); @@ -157,6 +159,45 @@ describe('workflow > triggers > schedule > date field mode', () => { expect(d0).toBe(startTime.getTime()); }); + it('starts on post.createdAt with offset in hours', async () => { + const workflow = await WorkflowModel.create({ + enabled: true, + type: 'schedule', + config: { + mode: 1, + collection: 'posts', + startsOn: { + field: 'createdAt', + offset: 1, + unit: 3600_000, + }, + repeat: 3000, + limit: 2, + }, + }); + + const now = await sleepToEvenSecond(); + + const before = new Date(); + before.setHours(before.getHours() - 2); + before.setSeconds(before.getSeconds() + 2); + + const post = await PostRepo.create({ values: { title: 't1', createdAt: before } }); + + await sleep(3000); + const e1s = await workflow.getExecutions(); + expect(e1s.length).toBe(1); + + await sleep(3000); + const e2s = await workflow.getExecutions(); + expect(e2s.length).toBe(2); + expect(e2s[0].context.data.id).toBe(post.id); + + // const triggerTime = new Date(post.createdAt.getTime() + 2000); + // triggerTime.setMilliseconds(0); + // expect(e2s[0].context.date).toBe(triggerTime.toISOString()); + }); + it('starts on post.createdAt and repeat by cron', async () => { const workflow = await WorkflowModel.create({ enabled: true, @@ -270,6 +311,7 @@ describe('workflow > triggers > schedule > date field mode', () => { endsOn: { field: 'createdAt', offset: 3, + unit: 1000, }, }, }); @@ -364,6 +406,7 @@ describe('workflow > triggers > schedule > date field mode', () => { endsOn: { field: 'createdAt', offset: 3, + unit: 1000, }, }, }); @@ -395,6 +438,7 @@ describe('workflow > triggers > schedule > date field mode', () => { startsOn: { field: 'createdAt', offset: 2, + unit: 1000, }, appends: ['category'], }, @@ -423,6 +467,7 @@ describe('workflow > triggers > schedule > date field mode', () => { endsOn: { field: 'createdAt', offset: 3, + unit: 1000, }, }, }); @@ -509,4 +554,41 @@ describe('workflow > triggers > schedule > date field mode', () => { expect(e2s.length).toBe(1); }); }); + + describe('record', () => { + it('record deleted after first triggered', async () => { + const workflow = await WorkflowModel.create({ + enabled: true, + type: 'schedule', + config: { + mode: 1, + collection: 'posts', + startsOn: { + field: 'createdAt', + }, + repeat: 1000, + }, + }); + + const now = await sleepToEvenSecond(); + + const p1 = await PostRepo.create({ values: { title: 't1' } }); + + await sleep(1300); + + const e1s = await workflow.getExecutions({ order: [['id', 'ASC']] }); + expect(e1s.length).toBe(1); + expect(e1s[0].context.data.id).toBe(p1.id); + const triggerTime = new Date(p1.createdAt); + triggerTime.setMilliseconds(0); + expect(e1s[0].context.date).toBe(triggerTime.toISOString()); + + await p1.destroy(); + + await sleep(1500); + + const e2s = await workflow.getExecutions({ order: [['id', 'ASC']] }); + expect(e2s.length).toBe(1); + }); + }); }); diff --git a/packages/plugins/@nocobase/plugin-workflow/src/server/triggers/ScheduleTrigger/DateFieldScheduleTrigger.ts b/packages/plugins/@nocobase/plugin-workflow/src/server/triggers/ScheduleTrigger/DateFieldScheduleTrigger.ts index 2be5136bb6b..28b8297bf4d 100644 --- a/packages/plugins/@nocobase/plugin-workflow/src/server/triggers/ScheduleTrigger/DateFieldScheduleTrigger.ts +++ b/packages/plugins/@nocobase/plugin-workflow/src/server/triggers/ScheduleTrigger/DateFieldScheduleTrigger.ts @@ -7,7 +7,7 @@ * For more information, please refer to: https://www.nocobase.com/agreement. */ -import { fn, literal, Op, Transactionable, where } from '@nocobase/database'; +import { fn, literal, Model, Op, Transactionable, where } from '@nocobase/database'; import parser from 'cron-parser'; import type Plugin from '../../Plugin'; import type { WorkflowModel } from '../../types'; @@ -36,7 +36,7 @@ export interface ScheduleTriggerConfig { endsOn?: string | ScheduleOnField; } -function getOnTimestampWithOffset({ field, offset = 0, unit = 1000 }: ScheduleOnField, now: Date) { +function getOnTimestampWithOffset({ field, offset = 0, unit = 86400000 }: ScheduleOnField, now: Date) { if (!field) { return null; } @@ -56,7 +56,7 @@ function getDataOptionTime(record, on, dir = 1) { return time ? time : null; } case 'object': { - const { field, offset = 0, unit = 1000 } = on; + const { field, offset = 0, unit = 86400000 } = on; if (!field || !record.get(field)) { return null; } @@ -206,17 +206,17 @@ export default class DateFieldScheduleTrigger { if (repeat) { // when repeat is number, means repeat after startsOn - // (now - startsOn) % repeat <= cacheCycle if (typeof repeat === 'number') { const tsFn = DialectTimestampFnMap[db.options.dialect]; if (repeat > range && tsFn) { + const offsetSeconds = Math.round(((startsOn.offset || 0) * (startsOn.unit || 1000)) / 1000); + const repeatSeconds = Math.round(repeat / 1000); + const nowSeconds = Math.round(timestamp / 1000); const { field } = model.getAttributes()[startsOn.field]; - const modExp = fn( - 'MOD', - literal( - `${Math.round(timestamp / 1000)} - ${tsFn(db.sequelize.getQueryInterface().quoteIdentifiers(field))}`, - ), - Math.round(repeat / 1000), + const modExp = literal( + `MOD(MOD(${tsFn( + db.sequelize.getQueryInterface().quoteIdentifiers(field), + )} + ${offsetSeconds} - ${nowSeconds}, ${repeatSeconds}) + ${repeatSeconds}, ${repeatSeconds})`, ); conditions.push(where(modExp, { [Op.lt]: Math.round(range / 1000) })); } @@ -250,6 +250,7 @@ export default class DateFieldScheduleTrigger { }); } this.workflow.getLogger(id).debug(`[Schedule on date field] conditions: `, { conditions }); + return model.findAll({ where: { [Op.and]: conditions, @@ -257,7 +258,7 @@ export default class DateFieldScheduleTrigger { }); } - getRecordNextTime(workflow: WorkflowModel, record, nextSecond = false) { + getRecordNextTime(workflow: WorkflowModel, record: Model, nextSecond = false) { const { config: { startsOn, endsOn, repeat, limit }, stats, @@ -265,6 +266,7 @@ export default class DateFieldScheduleTrigger { if (limit && stats.executed >= limit) { return null; } + const logger = this.workflow.getLogger(workflow.id); const range = this.cacheCycle; const now = new Date(); now.setMilliseconds(nextSecond ? 1000 : 0); @@ -273,43 +275,56 @@ export default class DateFieldScheduleTrigger { const endTime = getDataOptionTime(record, endsOn); let nextTime = null; if (!startTime) { + logger.debug(`[Schedule on date field] getNextTime: startsOn not configured`); return null; } if (startTime > timestamp + range) { + logger.debug(`[Schedule on date field] getNextTime: startsOn is out of caching window`); return null; } if (startTime >= timestamp) { - return !endTime || (endTime >= startTime && endTime < timestamp + range) ? startTime : null; + if (!endTime || startTime <= endTime) { + return startTime; + } + logger.debug(`[Schedule on date field] getNextTime: endsOn is before startsOn or out of caching window`); + return null; } else { if (!repeat) { + logger.debug( + `[Schedule on date field] getNextTime: startsOn is before current time and repeat is not configured`, + ); return null; } } if (typeof repeat === 'number') { - const nextRepeatTime = ((startTime - timestamp) % repeat) + repeat; - if (nextRepeatTime > range) { - return null; - } - if (endTime && endTime < timestamp + nextRepeatTime) { - return null; - } - nextTime = timestamp + nextRepeatTime; - } else if (typeof repeat === 'string') { - nextTime = getCronNextTime(repeat, now); + nextTime = timestamp + repeat - ((timestamp - startTime) % repeat); if (nextTime - timestamp > range) { + logger.debug(`[Schedule on date field] getNextTime: nextTime (${nextTime}) is out of caching window`); return null; } if (endTime && endTime < nextTime) { + logger.debug(`[Schedule on date field] getNextTime: nextTime is after endsOn`); + return null; + } + } else if (typeof repeat === 'string') { + nextTime = getCronNextTime(repeat, now); + if (nextTime - timestamp > range) { + logger.debug(`[Schedule on date field] getNextTime: nextTime (${nextTime}) is out of caching window`); + return null; + } + if (endTime && endTime < nextTime) { + logger.debug(`[Schedule on date field] getNextTime: nextTime is after endsOn`); return null; } } if (endTime && endTime <= timestamp) { + logger.debug(`[Schedule on date field] getNextTime: nextTime is after endsOn`); return null; } return nextTime; } - schedule(workflow: WorkflowModel, record, nextTime, toggle = true, options = {}) { + schedule(workflow: WorkflowModel, record: Model, nextTime: number, toggle = true, options = {}) { const [dataSourceName, collectionName] = parseCollectionName(workflow.config.collection); const { filterTargetKey } = this.workflow.app.dataSourceManager.dataSources .get(dataSourceName) @@ -336,17 +351,23 @@ export default class DateFieldScheduleTrigger { } } - async trigger(workflow: WorkflowModel, record, nextTime, { transaction }: Transactionable = {}) { + async trigger(workflow: WorkflowModel, record: Model, nextTime: number, { transaction }: Transactionable = {}) { const [dataSourceName, collectionName] = parseCollectionName(workflow.config.collection); const { repository, filterTargetKey } = this.workflow.app.dataSourceManager.dataSources .get(dataSourceName) .collectionManager.getCollection(collectionName); const recordPk = record.get(filterTargetKey); - const data = await repository.findOne({ + const data = (await repository.findOne({ filterByTk: recordPk, appends: workflow.config.appends, transaction, - }); + })) as Model; + if (!data) { + this.workflow + .getLogger(workflow.id) + .warn(`[Schedule on date field] record (${recordPk}) not exists, will not trigger`); + return; + } const eventKey = `${workflow.id}:${recordPk}@${nextTime}`; this.cache.delete(eventKey); // NOTE: data.toJSON() will cause erorr @@ -391,8 +412,9 @@ export default class DateFieldScheduleTrigger { return; } - const listener = async (data, { transaction }) => { + const listener = async (data: Model, { transaction }) => { const nextTime = this.getRecordNextTime(workflow, data); + this.workflow.getLogger().debug(`[Schedule on date field] record saved, nextTime: ${nextTime}`); return this.schedule(workflow, data, nextTime, Boolean(nextTime), { transaction }); };