Skip to content
Closed

wip #775

Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion examples/schedule.js
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ async function schedule () {

await boss.createQueue(queue)

await boss.schedule(queue, '*/2 * * * *', { arg1: 'schedule me' })
await boss.schedule(queue, '*/5 * * * * *', { arg1: 'schedule me' })

await boss.work(queue, async ([job]) => {
console.log(`received job ${job.id} with data ${JSON.stringify(job.data)} on ${new Date().toISOString()}`)
Expand Down
75 changes: 64 additions & 11 deletions src/timekeeper.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,17 @@ const WARNING_TYPES = {
CLOCK_SKEW: 'clock_skew'
} as const

function isSubMinuteCron (cron: string): boolean {
return cron.trim().split(/\s+/).length === 6
}

function getCronIntervalSeconds (cron: string, tz: string): number {
const interval = CronExpressionParser.parse(cron, { tz, strict: false })
const prev = interval.prev()
const next = interval.next()
return Math.round((next.getTime() - prev.getTime()) / 1000)
}

class Timekeeper extends EventEmitter implements types.EventsMixin {
db: types.IDatabase
config: types.ResolvedConstructorOptions
Expand All @@ -38,6 +49,8 @@ class Timekeeper extends EventEmitter implements types.EventsMixin {
private skewMonitorInterval: NodeJS.Timeout | null | undefined
private timekeeping: boolean | undefined
private _checkingSkew = false
private schedules: types.Schedule[] = []
private secondTickInterval: NodeJS.Timeout | null | undefined

clockSkew = 0
events = EVENTS
Expand Down Expand Up @@ -78,10 +91,12 @@ class Timekeeper extends EventEmitter implements types.EventsMixin {

await this.manager.work<types.Request>(QUEUES.SEND_IT, options, (jobs) => this.onSendIt(jobs))

await this.refreshScheduleCache()
setImmediate(() => this.onCron())

this.cronMonitorInterval = setInterval(async () => await this.onCron(), this.config.cronMonitorIntervalSeconds! * 1000)
this.skewMonitorInterval = setInterval(async () => await this.cacheClockSkew(), this.config.clockMonitorIntervalSeconds! * 1000)
this.secondTickInterval = setInterval(async () => await this.onSecondTick(), 1000)
}

async stop () {
Expand All @@ -103,6 +118,11 @@ class Timekeeper extends EventEmitter implements types.EventsMixin {
this.cronMonitorInterval = null
}

if (this.secondTickInterval) {
clearInterval(this.secondTickInterval)
this.secondTickInterval = null
}

while (this.timekeeping || this._checkingSkew) {
await delay(10)
}
Expand Down Expand Up @@ -134,10 +154,10 @@ class Timekeeper extends EventEmitter implements types.EventsMixin {

if (skewSeconds >= 60 || this.config.__test__force_clock_skew_warning) {
await emitAndPersistWarning(
this.warningContext,
WARNING_TYPES.CLOCK_SKEW,
WARNINGS.CLOCK_SKEW.message,
{ seconds: skewSeconds, direction: skew > 0 ? 'slower' : 'faster' }
this.warningContext,
WARNING_TYPES.CLOCK_SKEW,
WARNINGS.CLOCK_SKEW.message,
{ seconds: skewSeconds, direction: skew > 0 ? 'slower' : 'faster' }
)
}
} catch (err) {
Expand All @@ -158,6 +178,8 @@ class Timekeeper extends EventEmitter implements types.EventsMixin {

this.timekeeping = true

await this.refreshScheduleCache()

const sql = plans.trySetCronTime(this.config.schema, this.config.cronMonitorIntervalSeconds)

if (!this.stopped) {
Expand All @@ -175,18 +197,39 @@ class Timekeeper extends EventEmitter implements types.EventsMixin {
}

async cron () {
const schedules = await this.getSchedules()

const scheduled = schedules
.filter(i => this.shouldSendIt(i.cron, i.timezone))
.map(({ name, key, data, options }): types.JobInsert => ({ data: { name, data, options }, singletonKey: `${name}__${key}`, singletonSeconds: 60 }))
const scheduled = this.schedules
.filter(i => !isSubMinuteCron(i.cron) && this.shouldSendIt(i.cron, i.timezone))
.map(({ name, key, data, options }): types.JobInsert => ({ data: { name, data, options }, singletonKey: `${name}__${key}`, singletonSeconds: 60 }))

if (scheduled.length > 0 && !this.stopped) {
await this.manager.insert(QUEUES.SEND_IT, scheduled)
}
}

shouldSendIt (cron: string, tz: string) {
async onSecondTick () {
if (this.stopped) return

try {
const scheduled = this.schedules
.filter(i => isSubMinuteCron(i.cron) && this.shouldSendIt(i.cron, i.timezone, 1.5))
.map(({ name, key, data, options, cron, timezone }): types.JobInsert => {
const intervalSeconds = getCronIntervalSeconds(cron, timezone)
return {
data: { name, data, options },
singletonKey: `${name}__${key}`,
...(intervalSeconds > 1 ? { singletonSeconds: Math.min(intervalSeconds, 60) } : {})
}
})

if (scheduled.length > 0) {
await this.manager.insert(QUEUES.SEND_IT, scheduled)
}
} catch (err) {
this.emit(this.events.error, err)
}
}

shouldSendIt (cron: string, tz: string, windowSeconds = 60) {
const interval = CronExpressionParser.parse(cron, { tz, strict: false })

const prevTime = interval.prev()
Expand All @@ -195,7 +238,15 @@ class Timekeeper extends EventEmitter implements types.EventsMixin {

const prevDiff = (databaseTime - prevTime.getTime()) / 1000

return prevDiff < 60
return prevDiff < windowSeconds
}

async refreshScheduleCache (): Promise<void> {
try {
this.schedules = await this.getSchedules()
} catch (err) {
this.emit(this.events.error, err)
}
}

private async onSendIt (jobs: types.Job<types.Request>[]): Promise<void> {
Expand Down Expand Up @@ -227,6 +278,7 @@ class Timekeeper extends EventEmitter implements types.EventsMixin {
try {
const sql = plans.schedule(this.config.schema)
await this.db.executeSql(sql, [name, key, cron, tz, data, options])
await this.refreshScheduleCache()
} catch (err: any) {
if (err.message.includes('foreign key')) {
err.message = `Queue ${name} not found`
Expand All @@ -239,6 +291,7 @@ class Timekeeper extends EventEmitter implements types.EventsMixin {
async unschedule (name: string, key = ''): Promise<void> {
const sql = plans.unschedule(this.config.schema)
await this.db.executeSql(sql, [name, key])
await this.refreshScheduleCache()
}
}

Expand Down
66 changes: 66 additions & 0 deletions test/scheduleTest.ts
Original file line number Diff line number Diff line change
Expand Up @@ -365,6 +365,72 @@ describe('schedule', function () {
expect(schedules[0].cron).toBe('0 1 * * *')
})

it('should send jobs based on every one second expression', async function () {
const config = {
...ctx.bossConfig,
cronWorkerIntervalSeconds: 1,
schedule: true
}

ctx.boss = await helper.start(config)

await ctx.boss.schedule(ctx.schema, '* * * * * *')

let numJobs = 0
await ctx.boss.work(ctx.schema, async () => { numJobs++ })

await delay(3000)

expect(numJobs).toBeGreaterThan(1)
})

it('should send jobs based on every five seconds expression', async function () {
const config = {
...ctx.bossConfig,
cronWorkerIntervalSeconds: 1,
schedule: true
}

ctx.boss = await helper.start(config)

await ctx.boss.schedule(ctx.schema, '*/5 * * * * *')

let numJobs = 0
await ctx.boss.work(ctx.schema, async () => { numJobs++ })

await delay(7000)

expect(numJobs).toBeGreaterThanOrEqual(1)
})

it('in case of restart, jobs should not be overscheduled', async function () {
const config = {
...ctx.bossConfig,
cronWorkerIntervalSeconds: 1,
schedule: true
}

ctx.boss = await helper.start(config)

await ctx.boss.schedule(ctx.schema, '*/5 * * * * *')

let numJobs = 0
await ctx.boss.work(ctx.schema, async () => { numJobs++ })

await delay(6000)

await ctx.boss.stop({ graceful: false })

ctx.boss = await helper.start(config)
await ctx.boss.work(ctx.schema, async () => { numJobs++ })

await delay(1500)

// 7.5s total: */5 * * * * * fires at most 2 times in any 7.5s window
// a restart that re-fires the same occurrence would push the count to 3+
expect(numJobs).toBeLessThanOrEqual(2)
})

it('should get schedules filtered by a queue name and key', async function () {
const config = {
...ctx.bossConfig
Expand Down