@ -17,8 +17,18 @@ import {
import { and , eq , inArray , lt , sql } from 'drizzle-orm'
import type { PoolClient } from 'pg'
/** What the scheduler hands an in-process task besides its payload. */
export interface TaskContext {
/ * *
* Aborted when the task has run for ` scheduler.taskTimeout ` . A task cannot be stopped from outside
* the way a worker thread can , so this is a request : one that does its work in steps should check
* it between them and stop , and one that ignores it is abandoned — see ` executeInProcess ` .
* /
signal : AbortSignal
}
/** An in-process task, loaded from `tasks/simple/`. */
export type SimpleTask = ( payload? : any ) = > Promise < void > | void
export type SimpleTask = ( payload : any , context : TaskContext ) = > Promise < void > | void
/** Fallback for `scheduler.taskTimeout`, in seconds, when nothing is configured. */
const DEFAULT_TASK_TIMEOUT = 300
@ -35,6 +45,25 @@ const DEFAULT_STALE_JOB_TIMEOUT = 3600
* /
const TASK_TIMEOUT_GRACE = 5000
/ * *
* How long an in - process task that has been asked to stop gets to do so before it is abandoned .
*
* Far longer than the worker ' s grace , because stopping here is cooperative : a task stops at the end
* of the step it is in — the page it is rendering , the target it is syncing — rather than the moment
* the signal fires , and one of those steps may itself take half a minute .
* /
const IN_PROCESS_STOP_GRACE = 60 _000
/ * *
* In - process tasks that timed out without stopping , and when they were abandoned .
*
* Such a task is still running and nothing can end it , so no second copy of it is started on this
* instance until it settles . A task hangs on something — a lock , a connection , a dead browser — and a
* second copy almost always hangs on the same thing , each one taking a slot for the length of the
* timeout ; a task scheduled every minute would fill the pool within a few of them .
* /
const abandonedTasks = new Map < string , number > ( )
/ * *
* Sends the scheduler ' s cross - instance notifications , one at a time .
*
@ -267,6 +296,67 @@ export default {
}
} ,
/ * *
* Run a job in this process , and stop waiting for it if it does not come back .
*
* The counterpart of ` executeOnWorker ` , and without it one task that never settles is fatal rather
* than slow : its slot is never given back , its history row says ` active ` for ever , and a task that
* is scheduled keeps queuing copies that wait on the same thing .
*
* There is no thread to abort , so the two ceilings mean something different here . At the timeout the
* task ' s signal is aborted , which a task that works in steps answers by stopping at the end of the
* current one . If it has not settled a grace period later it is abandoned : recorded as failed like a
* worker that never answered , and left running , since nothing in JavaScript can stop a promise .
* /
async executeInProcess ( job : { task : string ; payload? : any } ) : Promise < void > {
const abandonedAt = abandonedTasks . get ( job . task )
if ( abandonedAt !== undefined ) {
throw new Error (
` An earlier run of ${ job . task } timed out ${ Math . round ( ( Temporal . Now . instant ( ) . epochMilliseconds - abandonedAt ) / 1000 ) } s ago and is still running, so another is not being started alongside it. `
)
}
const timeoutSeconds = WIKI . config . scheduler . taskTimeout ? ? DEFAULT_TASK_TIMEOUT
const controller = new AbortController ( )
const stopTimer = setTimeout ( ( ) = > {
controller . abort (
new Error ( ` The task did not finish within ${ timeoutSeconds } s and was asked to stop. ` )
)
} , timeoutSeconds * 1000 )
let giveUpTimer : NodeJS.Timeout | undefined
const run = Promise . resolve ( ) . then ( ( ) = >
this . tasks ! [ job . task ] ( job . payload , { signal : controller.signal } )
)
const abandon = ( ) = > {
abandonedTasks . set ( job . task , Temporal . Now . instant ( ) . epochMilliseconds )
// -> Also what keeps a late rejection from surfacing as an unhandled one
void run
. then (
( ) = > WIKI . logger . warn ( ` The abandoned run of ${ job . task } has finished after all. ` ) ,
( err : any ) = > WIKI . logger . warn ( ` The abandoned run of ${ job . task } failed: ${ err . message } ` )
)
. finally ( ( ) = > abandonedTasks . delete ( job . task ) )
return new Error (
` The task did not finish within ${ timeoutSeconds } s, nor stop within ${ IN_PROCESS_STOP_GRACE / 1000 } s of being asked to. It is still running and cannot be stopped from here; no other run of it will start on this instance until it ends. `
)
}
try {
await Promise . race ( [
run ,
new Promise < never > ( ( _resolve , reject ) = > {
giveUpTimer = setTimeout (
( ) = > reject ( abandon ( ) ) ,
timeoutSeconds * 1000 + IN_PROCESS_STOP_GRACE
)
} )
] )
} finally {
clearTimeout ( stopTimer )
clearTimeout ( giveUpTimer )
}
} ,
/ * *
* Take a batch of due jobs and run them .
*
@ -359,7 +449,7 @@ export default {
if ( job . useWorker ) {
await this . executeOnWorker ( job )
} else {
await this . tasks! [ job . task ] ( job . payload )
await this . executeInProcess( job )
}
await WIKI . db
. update ( jobHistoryTable )