feat(events): the runner (Phase 2)
`utils/eventRunner.js`, the eighth poller, wired into server.js beside engagementWorker. Its tick reclaims stale leases, sweeps occurrences past their grace window into `missed`, advances each due run through its phases, and drains that phase's steps in `seq` order. The three core actions from Phase 1 get real bodies, so a published event started from the existing run route now announces, waits and completes on its own. No routes are added: a runner has no surface, and the live controls stay Phase 3's. Four things the org lead settled (2026-09-02): a parked step is `running` with a NULL lease; `await: 'human'` and `holdFor` are ordinary success-envelope members rather than special cases keyed on an action id; a run whose concurrency key is held stays `scheduled` and lets its grace window decide; and `n` in §L's `retry(n)` is a runner constant. Co-Authored-By: Claude <noreply@anthropic.com>
This commit is contained in:
@@ -13,28 +13,20 @@
|
||||
// a human to go and do something. A deployment with no game module installed has
|
||||
// a working event system made of exactly these.
|
||||
//
|
||||
// **Nothing here dispatches yet.** Phase 1 builds the registry, the id grammar,
|
||||
// the risk classes and the param validation; Phase 2 builds `utils/eventRunner.js`
|
||||
// and is what calls `perform()`. The bodies below therefore answer with the
|
||||
// envelope §F defines for a refusal — and specifically NOT with `{ ok: true }`,
|
||||
// which is the one wrong answer a placeholder can give: `ok: true` on an action
|
||||
// that did nothing is a recorded world change that did not occur, which is the
|
||||
// exact mistake the envelope's failure default exists to prevent. `retry: false`
|
||||
// because a missing runner is not a transient condition.
|
||||
// **Phase 2 gave all three real bodies**, and between them they exercise every
|
||||
// shape §F's envelope can take: `core.announce` does work and finishes,
|
||||
// `core.wait` finishes while deferring what follows it, and `core.cue` succeeds
|
||||
// without finishing at all. The runner learns nothing about any of them by id —
|
||||
// each says what it needs in the envelope, through the same two members Phase 7
|
||||
// hands to a module.
|
||||
//
|
||||
// **This file must not touch the database.** It is required from `registerCore()`,
|
||||
// which runs under `routeManifest.js` and `swagger.js` against a dead pool
|
||||
// (MODULE_API.md §2.2). It is pure data plus three functions that are not called.
|
||||
// (MODULE_API.md §2.2). Nothing below runs at require time; the announce leg is
|
||||
// looked up inside `perform()`, per call, which is also what makes a leg
|
||||
// registered by a module that booted later reachable at all.
|
||||
|
||||
// A factory rather than one shared function, because `perform`'s argument is
|
||||
// §F's dispatch envelope — `{ runId, stepId, idempotencyKey, scope, params,
|
||||
// actor, verify }` — and it does not carry the action's own id. Closing over it
|
||||
// is what lets the refusal name which action refused.
|
||||
const notWiredYet = (actionId) => async () => ({
|
||||
ok: false,
|
||||
retry: false,
|
||||
error: `${actionId} is declared in Phase 1 and dispatched from Phase 2`,
|
||||
})
|
||||
const registries = require('../modules/registries')
|
||||
|
||||
const ACTIONS = [
|
||||
{
|
||||
@@ -81,7 +73,52 @@ const ACTIONS = [
|
||||
},
|
||||
],
|
||||
|
||||
perform: notWiredYet('core.announce'),
|
||||
/**
|
||||
* Publish through the announce leg the step names.
|
||||
*
|
||||
* **The legs are reused rather than reimplemented** (§J, "reuse the legs"):
|
||||
* `discord` is core's and `towncrier` is module-uo's, both already registered,
|
||||
* both already carrying a `classify()` that knows what their transport's
|
||||
* failures mean. An event announcement that went out by some other path would
|
||||
* be a second delivery mechanism with its own bugs.
|
||||
*
|
||||
* A leg's `dispatch()` takes a POST — that is the shape the news path gave it
|
||||
* — so an event announcement is presented as one. `excerpt` is the body
|
||||
* because it is the field every leg renders as prose, and `image_url` is null
|
||||
* because an event announcement has no article behind it to illustrate.
|
||||
* Widening the leg contract to carry a second payload shape is a
|
||||
* MODULE_API change, and Phase 7 is where those are made.
|
||||
*
|
||||
* The leg id is checked HERE rather than at authoring time, and that is not
|
||||
* laxness: legs are registered by modules, and a spec is validated in a
|
||||
* process that may have booted before the module that owns the leg.
|
||||
*/
|
||||
async perform({ params, verify }) {
|
||||
const registered = registries.announceLeg(params.leg)
|
||||
if (!registered) {
|
||||
// Terminal, not transient. A leg nobody registers will not appear
|
||||
// between two attempts sixty seconds apart, and the honest cause — a
|
||||
// module removed, or a typo the authoring form could not catch — is a
|
||||
// thing a human fixes.
|
||||
return { ok: false, retry: false, error: `no module registers the announce leg "${params.leg}"` }
|
||||
}
|
||||
// A dry run reports what it WOULD do and sends nothing (§I). Answering
|
||||
// before the dispatch rather than inside the leg is what keeps that true
|
||||
// for legs written by people who never read this file.
|
||||
if (verify) return { ok: true }
|
||||
|
||||
const result = await registered.dispatch({
|
||||
title: params.title || null,
|
||||
excerpt: params.body,
|
||||
image_url: null,
|
||||
})
|
||||
// The leg's own classification, not a second opinion. `retry` vs
|
||||
// `terminal` for a Discord webhook is a judgement `discordAnnounce.classify`
|
||||
// already makes, and making it twice is how the two drift.
|
||||
const { outcome, error } = registered.classify(result)
|
||||
if (outcome === 'done') return { ok: true }
|
||||
return { ok: false, retry: outcome === 'retry', error: error || `announce leg "${params.leg}" refused` }
|
||||
},
|
||||
},
|
||||
|
||||
{
|
||||
@@ -105,11 +142,18 @@ const ACTIONS = [
|
||||
},
|
||||
],
|
||||
|
||||
// A wait is a genuine no-op at dispatch, and it will stay one: the delay is
|
||||
// the NEXT step's `due_at`, which the runner owns, not something this
|
||||
// function sleeps through. A `perform` that slept would hold a step's claim
|
||||
// for the duration and turn a five-minute pause into a five-minute lease.
|
||||
perform: notWiredYet('core.wait'),
|
||||
// A wait is a genuine no-op at dispatch, and it stayed one: the delay is the
|
||||
// NEXT step's `due_at`, which the runner owns, not something this function
|
||||
// sleeps through. A `perform` that slept would hold a step's claim for the
|
||||
// duration and turn a five-minute pause into a five-minute lease — and the
|
||||
// reclaim would then re-dispatch it, so a long enough wait would never end.
|
||||
//
|
||||
// `holdFor` is an ordinary envelope member (org lead, 2026-09-02), which is
|
||||
// why the runner can honour this without knowing what `core.wait` is.
|
||||
async perform({ params, verify }) {
|
||||
if (verify) return { ok: true }
|
||||
return { ok: true, holdFor: params.seconds }
|
||||
},
|
||||
},
|
||||
|
||||
{
|
||||
@@ -142,11 +186,26 @@ const ACTIONS = [
|
||||
},
|
||||
],
|
||||
|
||||
// Phase 2 gives this its parking semantics — a cue step does not complete
|
||||
// when `perform` answers, it completes when a human presses confirm, and the
|
||||
// control that does so is Phase 3's. Both of those are what make this the
|
||||
// one action whose runtime shape is deliberately not decided here.
|
||||
perform: notWiredYet('core.cue'),
|
||||
/**
|
||||
* Post the instruction and PARK. The step does not complete here.
|
||||
*
|
||||
* `await: 'human'` is the envelope member that says so (org lead,
|
||||
* 2026-09-02), and the runner's answer to it is to leave the step `running`
|
||||
* with a NULL lease — genuinely in flight, nothing holding it, so the stale
|
||||
* reclaim passes it by and a cue posted on Friday is still waiting on Monday.
|
||||
* The step ends when someone presses confirm, which is Phase 3's control.
|
||||
*
|
||||
* **Nothing is delivered from here in Phase 2, and that is visible rather
|
||||
* than pretended.** The instruction is carried by the step's own params and
|
||||
* shown on the run console; routing it to Discord or to a staff inbox is
|
||||
* Phase 10's integration work, through the engagement triggers that own every
|
||||
* other notification on this platform. An action that grew its own delivery
|
||||
* path would be the second one.
|
||||
*/
|
||||
async perform({ verify }) {
|
||||
if (verify) return { ok: true }
|
||||
return { ok: true, await: 'human' }
|
||||
},
|
||||
},
|
||||
]
|
||||
|
||||
|
||||
166
server/src/events/dispatch.js
Normal file
166
server/src/events/dispatch.js
Normal file
@@ -0,0 +1,166 @@
|
||||
// ── Dispatching one step to one action ─────────────────────────────────────
|
||||
//
|
||||
// EVENTS.md §F. This is the boundary between the runner and code core did not
|
||||
// write, and it exists as its own file because it has exactly one job: call
|
||||
// `perform()` and turn whatever comes back — an envelope, a lie, a throw, a
|
||||
// promise that never settles — into one of four classifications the runner knows
|
||||
// how to act on.
|
||||
//
|
||||
// **§F's load-bearing rule, and the reason none of this is inlined into the
|
||||
// runner: no shape a failure can take may read as success.** A rejected promise,
|
||||
// a throw, a timeout, a non-object and a missing `ok` are all
|
||||
// `{ ok: false, retry: true }`. That is the inverse of `registerTeamProvider`'s
|
||||
// default, deliberately — a team provider that refuses leaves core showing what
|
||||
// it already had, because staleness is cheap, whereas an action that half-ran and
|
||||
// was recorded as `done` is a world change nothing will ever come back for.
|
||||
//
|
||||
// **The timeout is the module contract's, not this file's opinion.** Every action
|
||||
// declares `budgetMs` at registration and the registry bounds it there; here it
|
||||
// is enforced. Without it a module whose `perform()` awaits a socket that never
|
||||
// answers holds a step's claim until the lease expires, and the reclaim then
|
||||
// re-dispatches it — which is how one wedged sidecar becomes an infinite loop
|
||||
// rather than a failed step.
|
||||
|
||||
const registries = require('../modules/registries')
|
||||
const log = require('../utils/logger')('events')
|
||||
|
||||
// What a classification can be. `parked` is Phase 2's addition and it is the one
|
||||
// outcome that is neither terminal nor a retry: the action succeeded, and the
|
||||
// step is not finished, because something outside this system has to happen next.
|
||||
const OUTCOMES = ['done', 'parked', 'retry', 'terminal']
|
||||
|
||||
// The upper bound on `holdFor`, in seconds. A wait is a scheduling instruction,
|
||||
// not a lease, so this is generous — but it is bounded, because an action that
|
||||
// answers `holdFor: 1e9` would park the phase past the heat death of the shard
|
||||
// and the step that did it would look, in the console, exactly like one that
|
||||
// worked.
|
||||
const MAX_HOLD_SECONDS = 7 * 24 * 60 * 60
|
||||
|
||||
/**
|
||||
* Run `fn()` under a deadline.
|
||||
*
|
||||
* The loser of the race is not cancelled — JavaScript has no such thing, and a
|
||||
* `perform()` still awaiting a socket keeps awaiting it. What the deadline buys
|
||||
* is that the RUNNER stops waiting, which is the half that matters: the step is
|
||||
* classified, the claim is released, and the tick moves on. A late answer from
|
||||
* the abandoned call lands on a step that has already been written, and the
|
||||
* idempotency key is what makes the retry that follows safe on the game side.
|
||||
*/
|
||||
function withDeadline(fn, ms, actionId) {
|
||||
let timer = null
|
||||
const deadline = new Promise((resolve) => {
|
||||
timer = setTimeout(
|
||||
() => resolve({ __timedOut: true, error: `${actionId} exceeded its ${ms}ms budget` }),
|
||||
ms,
|
||||
)
|
||||
if (timer.unref) timer.unref()
|
||||
})
|
||||
return Promise.race([Promise.resolve().then(fn), deadline]).finally(() => {
|
||||
if (timer) clearTimeout(timer)
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
* Turn a raw `perform()` answer into `{ outcome, error?, holdSeconds?, resources? }`.
|
||||
*
|
||||
* Exported and pure, so the classification rules are testable without a registry,
|
||||
* a database or a clock — which matters because they are the rules that decide
|
||||
* whether a world change is recorded as having happened.
|
||||
*/
|
||||
function classify(result, actionId) {
|
||||
if (result && result.__timedOut) {
|
||||
// Transient by default: a timeout says nothing about whether the action ran.
|
||||
// That ambiguity is exactly what the idempotency key exists to resolve, and
|
||||
// resolving it on the game side is Phase 11's protocol work — until then a
|
||||
// retry is the honest choice and the risk class decides what happens when the
|
||||
// retries run out.
|
||||
return { outcome: 'retry', error: result.error }
|
||||
}
|
||||
if (result === null || typeof result !== 'object' || Array.isArray(result)) {
|
||||
return { outcome: 'retry', error: `${actionId} answered with no envelope` }
|
||||
}
|
||||
if (result.ok !== true) {
|
||||
// `retry` must be opted into. An action that means "this will never work"
|
||||
// says `retry: false`, and an envelope that forgot to say anything gets the
|
||||
// benefit of the doubt on the transient question but not on the success one.
|
||||
const retry = result.retry !== false
|
||||
return {
|
||||
outcome: retry ? 'retry' : 'terminal',
|
||||
error: result.error ? String(result.error) : `${actionId} refused`,
|
||||
}
|
||||
}
|
||||
|
||||
// ── The two success shapes that are not "finished" ──
|
||||
//
|
||||
// Both were settled by the org lead on 2026-09-02, and both are envelope
|
||||
// members rather than special cases keyed on an action id, so that the runner
|
||||
// never names a verb. `core.cue` and `core.wait` reach them through the same
|
||||
// door Phase 7 opens to a module's own long-running action.
|
||||
if (result.await === 'human') {
|
||||
return { outcome: 'parked', error: null, resources: result.resources || [] }
|
||||
}
|
||||
|
||||
let holdSeconds = 0
|
||||
if (result.holdFor !== undefined && result.holdFor !== null) {
|
||||
const n = Number(result.holdFor)
|
||||
if (!Number.isFinite(n) || n < 0) {
|
||||
return { outcome: 'terminal', error: `${actionId} answered a bad holdFor "${result.holdFor}"` }
|
||||
}
|
||||
holdSeconds = Math.min(Math.floor(n), MAX_HOLD_SECONDS)
|
||||
}
|
||||
|
||||
return { outcome: 'done', error: null, holdSeconds, resources: result.resources || [] }
|
||||
}
|
||||
|
||||
/**
|
||||
* Dispatch one step. Never throws.
|
||||
*
|
||||
* `verify` rides through to `perform()` unchanged (§I's dry run, Phase 6's
|
||||
* route): `verify === true` means validate and report, change nothing. It is
|
||||
* passed from here rather than being a separate code path so that the dry run
|
||||
* exercises the real dispatcher — a dry run down a second path is a dry run of
|
||||
* the second path.
|
||||
*/
|
||||
async function dispatchStep(step, { run, actor = null, verify = false } = {}) {
|
||||
const action = registries.eventAction(step.action_id)
|
||||
if (!action) {
|
||||
// §L, verbatim: "a step naming one fails terminal with the module named, and
|
||||
// the run degrades rather than claiming success. Never a silent skip." The
|
||||
// module was uninstalled or failed to boot between publish and now — publish
|
||||
// refuses a dormant step, so this cannot be an authoring mistake.
|
||||
return { outcome: 'terminal', error: `no module registers "${step.action_id}"`, dormant: true }
|
||||
}
|
||||
|
||||
const envelope = {
|
||||
runId: run.id,
|
||||
stepId: step.id,
|
||||
idempotencyKey: step.idempotency_key,
|
||||
scope: run.scope || '',
|
||||
params: step.params || {},
|
||||
actor,
|
||||
verify: Boolean(verify),
|
||||
}
|
||||
|
||||
let raw
|
||||
try {
|
||||
raw = await withDeadline(() => action.perform(envelope), action.budgetMs, action.id)
|
||||
} catch (err) {
|
||||
// A module should not throw, and if one does it is a transient failure rather
|
||||
// than a crashed tick — announceWorker's posture with its legs, and the
|
||||
// reason one bad module cannot stop every other run on the deployment.
|
||||
log.warn('event action threw', { action: action.id, run: run.id, step: step.id, message: err.message })
|
||||
return { outcome: 'retry', error: err.message }
|
||||
}
|
||||
|
||||
const classification = classify(raw, action.id)
|
||||
if (step.action_version && action.version !== step.action_version) {
|
||||
// Not a refusal: the step was authored against an older declaration and the
|
||||
// module has moved on. The editor is where that becomes a warning (§F); here
|
||||
// it is recorded, so a run that behaved oddly can be explained afterwards by
|
||||
// reading the log rather than by guessing.
|
||||
classification.actionVersionDrift = { authored: step.action_version, registered: action.version }
|
||||
}
|
||||
return classification
|
||||
}
|
||||
|
||||
module.exports = { dispatchStep, classify, withDeadline, OUTCOMES, MAX_HOLD_SECONDS }
|
||||
@@ -23,6 +23,16 @@ const KINDS = [
|
||||
'phase.entered', // a phase's steps were materialised
|
||||
'step.status', // a step transition, with the module's answer
|
||||
'note', // a human action taken from the admin surface
|
||||
// Phase 2's, all five of them answers to a question an operator asks out
|
||||
// loud. `run.blocked` in particular is the whole reason this table exists
|
||||
// rather than a server log line: "it did not start because run 37 holds
|
||||
// invasion:Yew" is a fact with two run ids in it, and it has to be
|
||||
// queryable from the run that did NOT start.
|
||||
'run.blocked', // an occurrence held off: another run has its concurrency key
|
||||
'run.health', // a health change, which is not a status change
|
||||
'step.retry', // a step failed transiently and will be attempted again
|
||||
'step.parked', // a step is waiting on a human and nothing is holding it
|
||||
'phase.completed', // every step of a phase reached a terminal status
|
||||
]
|
||||
|
||||
const hydrate = (row) => row && { ...row, detail: parseJson(row.detail, null) }
|
||||
@@ -64,4 +74,33 @@ async function write({ runId, stepId = null, kind, phase = null, detail = null }
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = { KINDS, listForRun, write }
|
||||
/**
|
||||
* Delete log lines belonging to runs that are both TERMINAL and older than
|
||||
* `before`, a bounded number at a time.
|
||||
*
|
||||
* The schema comment beside `idx_evlog_at` parked this sweep here, and it is the
|
||||
* rule Engagement Phase 14 arrived at applied to a second high-cardinality table:
|
||||
* **only terminal rows are eligible.** A run still in flight keeps every line it
|
||||
* has, however old — the log's whole job is answering "why didn't phase 3 start?"
|
||||
* about a run that is, right now, not starting phase 3, and a horizon that could
|
||||
* reach a live run would delete the answer while the question was still open.
|
||||
*
|
||||
* `LIMIT` makes one call a bounded amount of work rather than a table-sized
|
||||
* transaction; the timer runs again and takes the next slice. The join is on the
|
||||
* run's terminal status rather than on a precomputed id list so that a run which
|
||||
* reached a terminal state between the two would not be missed.
|
||||
*/
|
||||
const pruneTerminal = async (before, limit = 5000) => {
|
||||
const n = Math.min(Math.max(Number(limit) || 5000, 1), 50_000)
|
||||
const result = await query(
|
||||
`DELETE l FROM event_run_log l
|
||||
JOIN event_runs r ON r.id = l.run_id
|
||||
WHERE r.status IN ('completed','cancelled','failed','missed')
|
||||
AND COALESCE(r.ended_at, r.updated_at) < ?
|
||||
LIMIT ${n}`,
|
||||
[before],
|
||||
)
|
||||
return Number(result?.affectedRows || 0)
|
||||
}
|
||||
|
||||
module.exports = { KINDS, listForRun, write, pruneTerminal }
|
||||
|
||||
@@ -92,4 +92,201 @@ const statusCounts = async (runId) => {
|
||||
return Object.fromEntries(rows.map((r) => [r.status, Number(r.n)]))
|
||||
}
|
||||
|
||||
module.exports = { listForRun, getById, materialisePhase, statusCounts, idempotencyKey }
|
||||
// ── Phase 2: draining a step ───────────────────────────────────────────────
|
||||
//
|
||||
// The CAS claim, the lease, the attempt counter and the terminal writes. Phase 1
|
||||
// left all of it out rather than stubbing it, and this is where it lands.
|
||||
//
|
||||
// **Two rules govern everything below, and both were paid for once already.**
|
||||
//
|
||||
// 1. `attempts` is incremented by the CLAIM and by nothing else, and no recovery
|
||||
// path ever resets it. Engagement Phase 14's defect was a stale-row sweep that
|
||||
// returned rows to their start state: the attempt ceiling became unreachable,
|
||||
// so the row cycled forever, never reached a terminal status, and was
|
||||
// therefore never eligible for any retention sweep.
|
||||
// 2. A PARKED step is `running` with a NULL lease, and the reclaim only ever
|
||||
// touches a lease that is non-NULL and expired (the org lead's answer,
|
||||
// 2026-09-02). That is what lets a GM cue wait for a human overnight without a
|
||||
// sweep re-dispatching the instruction every fifteen minutes.
|
||||
|
||||
// A run's steps in authored order, for the phase the run is currently in.
|
||||
const listForPhase = async (runId, phase) =>
|
||||
(
|
||||
await query(
|
||||
'SELECT * FROM event_run_steps WHERE run_id = ? AND phase = ? ORDER BY seq, id',
|
||||
[runId, phase],
|
||||
)
|
||||
).map(hydrate)
|
||||
|
||||
/**
|
||||
* The next step of a phase that the runner may work on, or null.
|
||||
*
|
||||
* **Steps within a phase are strictly serial.** This returns the lowest-`seq`
|
||||
* step that is not terminal, and the runner does nothing with step N+1 until N
|
||||
* has finished — which is the only reading under which `core.wait` means anything
|
||||
* at all, and the only one under which a cue can gate what follows it.
|
||||
*
|
||||
* A parked or running step is returned too, so the caller can see that the phase
|
||||
* is occupied rather than concluding it is finished.
|
||||
*/
|
||||
const nextOpenStep = async (runId, phase) => {
|
||||
const [row] = await query(
|
||||
`SELECT * FROM event_run_steps
|
||||
WHERE run_id = ? AND phase = ?
|
||||
AND status IN ('pending','running')
|
||||
ORDER BY seq, id LIMIT 1`,
|
||||
[runId, phase],
|
||||
)
|
||||
return hydrate(row) || null
|
||||
}
|
||||
|
||||
/**
|
||||
* Take ownership of one pending step: the CAS `pending -> running`, plus a lease.
|
||||
*
|
||||
* `due_at` is honoured here rather than in the caller's filter so that the whole
|
||||
* decision — is it mine, is it due — is one statement the database arbitrates. A
|
||||
* NULL `due_at` is due now, which is what materialisation writes for every step
|
||||
* that is not sitting behind a `core.wait`.
|
||||
*/
|
||||
async function claim(id, owner, leaseUntil, now) {
|
||||
const result = await query(
|
||||
`UPDATE event_run_steps
|
||||
SET status = 'running', attempts = attempts + 1, claimed_by = ?, claim_expires_at = ?,
|
||||
started_at = COALESCE(started_at, NOW())
|
||||
WHERE id = ? AND status = 'pending' AND (due_at IS NULL OR due_at <= ?)`,
|
||||
[owner, leaseUntil, id, now],
|
||||
)
|
||||
return Number(result?.affectedRows || 0) === 1
|
||||
}
|
||||
|
||||
/**
|
||||
* Park a claimed step: it stays `running`, and its lease goes NULL.
|
||||
*
|
||||
* This is the whole mechanism behind `core.cue`. The step is genuinely in flight
|
||||
* — an instruction has been posted and nothing else in the phase may proceed —
|
||||
* but no process is holding it, so the reclaim must not take it back. A NULL
|
||||
* lease says exactly that, and `reclaimStale` below is written to agree.
|
||||
*/
|
||||
const park = (id, note) =>
|
||||
query(
|
||||
`UPDATE event_run_steps SET claim_expires_at = NULL, last_error = ?
|
||||
WHERE id = ? AND status = 'running'`,
|
||||
[note ? String(note).slice(0, 500) : null, id],
|
||||
)
|
||||
|
||||
/**
|
||||
* Release a claimed step back to `pending` for a later attempt.
|
||||
*
|
||||
* `attempts` is untouched — it was already incremented by the claim, which is the
|
||||
* only place that may. Backoff is flat rather than exponential for the reason the
|
||||
* outbox's is: `due_at` is also the event's own clock, and a doubling backoff
|
||||
* pushes a step arbitrarily far past the moment the event was about.
|
||||
*/
|
||||
const reschedule = (id, dueAt, error) =>
|
||||
query(
|
||||
`UPDATE event_run_steps
|
||||
SET status = 'pending', due_at = ?, claimed_by = NULL, claim_expires_at = NULL, last_error = ?
|
||||
WHERE id = ? AND status = 'running'`,
|
||||
[dueAt, error ? String(error).slice(0, 500) : null, id],
|
||||
)
|
||||
|
||||
/** A terminal outcome for one step: done, failed, skipped, refused or cancelled. */
|
||||
const finish = (id, status, error) =>
|
||||
query(
|
||||
`UPDATE event_run_steps
|
||||
SET status = ?, last_error = ?, finished_at = NOW(),
|
||||
claimed_by = NULL, claim_expires_at = NULL
|
||||
WHERE id = ? AND status = 'running'`,
|
||||
[status, error ? String(error).slice(0, 500) : null, id],
|
||||
)
|
||||
|
||||
/**
|
||||
* Delay the next not-yet-started step of a phase — what `core.wait` actually does.
|
||||
*
|
||||
* The wait step itself completes normally; the pause is the NEXT step's `due_at`,
|
||||
* owned by the runner. A `perform()` that slept would hold its claim for the
|
||||
* duration and turn a five-minute pause into a five-minute lease, which is the
|
||||
* one shape this must not have.
|
||||
*
|
||||
* Guarded on `status = 'pending'` and on the current `due_at` being sooner, so a
|
||||
* re-dispatch of a wait whose ack was lost cannot push the following step further
|
||||
* out a second time.
|
||||
*/
|
||||
const holdNext = async (runId, phase, afterSeq, dueAt) => {
|
||||
const result = await query(
|
||||
`UPDATE event_run_steps
|
||||
SET due_at = ?
|
||||
WHERE run_id = ? AND phase = ? AND seq > ? AND status = 'pending'
|
||||
AND (due_at IS NULL OR due_at < ?)
|
||||
ORDER BY seq LIMIT 1`,
|
||||
[dueAt, runId, phase, afterSeq, dueAt],
|
||||
)
|
||||
return Number(result?.affectedRows || 0) === 1
|
||||
}
|
||||
|
||||
/**
|
||||
* Recover steps whose claim outlived the process that took it.
|
||||
*
|
||||
* **`attempts` is not reset and the lease being NULL is not staleness.** The
|
||||
* first is Engagement Phase 14's rule; the second is what makes a parked cue
|
||||
* survive. A step that has already burned its attempts leaves `running` as
|
||||
* `failed` rather than being handed back, and in that order — a reclaim that ran
|
||||
* first would return it to `pending` and it would be retried forever.
|
||||
*/
|
||||
const reclaimStale = async (now, maxAttempts = 0) => {
|
||||
let failed = 0
|
||||
if (Number(maxAttempts) > 0) {
|
||||
const gaveUp = await query(
|
||||
`UPDATE event_run_steps
|
||||
SET status = 'failed', last_error = 'gave up after repeated interruptions',
|
||||
finished_at = NOW(), claimed_by = NULL, claim_expires_at = NULL
|
||||
WHERE status = 'running'
|
||||
AND claim_expires_at IS NOT NULL AND claim_expires_at < ?
|
||||
AND attempts >= ?`,
|
||||
[now, Math.floor(maxAttempts)],
|
||||
)
|
||||
failed = Number(gaveUp?.affectedRows || 0)
|
||||
}
|
||||
const reclaimed = await query(
|
||||
`UPDATE event_run_steps
|
||||
SET status = 'pending', claimed_by = NULL, claim_expires_at = NULL
|
||||
WHERE status = 'running' AND claim_expires_at IS NOT NULL AND claim_expires_at < ?`,
|
||||
[now],
|
||||
)
|
||||
return { failed, reclaimed: Number(reclaimed?.affectedRows || 0) }
|
||||
}
|
||||
|
||||
/**
|
||||
* Cancel every step of a run that has not started (Phase 3's cancel, and the
|
||||
* abort_run disposition).
|
||||
*
|
||||
* A `running` step is deliberately left alone, parked or not: nothing can recall
|
||||
* a command already sent, and a second writer on that row would race the process
|
||||
* that owns it (§L).
|
||||
*/
|
||||
const cancelPending = async (runId) => {
|
||||
const result = await query(
|
||||
`UPDATE event_run_steps
|
||||
SET status = 'cancelled', finished_at = NOW()
|
||||
WHERE run_id = ? AND status = 'pending'`,
|
||||
[runId],
|
||||
)
|
||||
return Number(result?.affectedRows || 0)
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
listForRun,
|
||||
listForPhase,
|
||||
getById,
|
||||
materialisePhase,
|
||||
statusCounts,
|
||||
idempotencyKey,
|
||||
nextOpenStep,
|
||||
claim,
|
||||
park,
|
||||
reschedule,
|
||||
finish,
|
||||
holdNext,
|
||||
reclaimStale,
|
||||
cancelPending,
|
||||
}
|
||||
|
||||
@@ -17,6 +17,12 @@ const { parseJson } = require('./eventJson')
|
||||
|
||||
const hydrate = (row) => row && { ...row, params: parseJson(row.params, null), rehearsal: Boolean(row.rehearsal) }
|
||||
|
||||
// The statuses a run never leaves. A transition INTO one of these stamps
|
||||
// `ended_at` and drops the claim, and only rows in one of them are eligible for
|
||||
// the log retention sweep -- Engagement Phase 14's rule, which is a bound at all
|
||||
// only if every path a run can take reaches one of them.
|
||||
const TERMINAL = ['completed', 'cancelled', 'failed', 'missed']
|
||||
|
||||
const SELECT_LIST = `
|
||||
SELECT r.*, d.title AS definition_title, d.slug AS definition_slug, v.version AS version_number
|
||||
FROM event_runs r
|
||||
@@ -102,4 +108,256 @@ const countActiveForDefinition = async (definitionId) => {
|
||||
return Number(row?.n || 0)
|
||||
}
|
||||
|
||||
module.exports = { list, getById, materialise, findOccurrence, countActiveForDefinition }
|
||||
// ── Phase 2: the claim, the transitions and the reclaim ────────────────────
|
||||
//
|
||||
// Everything below is the runner's, and none of it existed in Phase 1 for a
|
||||
// stated reason: a half-written claim is worse than no claim, because it reads
|
||||
// as protection. It is written here now, in full.
|
||||
//
|
||||
// **The division of labour with the unique index has not changed.** The index one
|
||||
// section up is what makes "one run per occurrence per scope" TRUE; the CAS below
|
||||
// decides only WHO advances an occurrence that already exists. Neither substitutes
|
||||
// for the other, and this deployment being single-instance (§N4) changes the test
|
||||
// rather than the design — the same two protections are what keep a tick that
|
||||
// overran into the next one from advancing a run twice.
|
||||
|
||||
/**
|
||||
* Runs the runner should look at this tick: due, and not yet terminal.
|
||||
*
|
||||
* It selects rather than claims — `claimStart` and `claimTick` below are one row
|
||||
* at a time — so two sweepers see the same candidates and then disagree,
|
||||
* harmlessly, about which of them owns each. `idx_evrun_due (status,
|
||||
* scheduled_for)` is this query.
|
||||
*
|
||||
* `paused` is absent from the status list on purpose. A paused run is waiting on
|
||||
* a human and must not be advanced by a tick; the only thing that moves it is
|
||||
* Phase 3's resume control.
|
||||
*/
|
||||
const findDue = async (now, limit = 50) => {
|
||||
const n = Math.min(Math.max(Number(limit) || 50, 1), 500)
|
||||
return (
|
||||
await query(
|
||||
`SELECT * FROM event_runs
|
||||
WHERE status IN ('scheduled','starting','running','ending')
|
||||
AND scheduled_for <= ?
|
||||
ORDER BY scheduled_for, id
|
||||
LIMIT ${n}`,
|
||||
[now],
|
||||
)
|
||||
).map(hydrate)
|
||||
}
|
||||
|
||||
/**
|
||||
* Take ownership of a run that has not started: the CAS `scheduled -> starting`.
|
||||
*
|
||||
* Verbatim the outbox claim the org lead settled over `SELECT ... FOR UPDATE
|
||||
* SKIP LOCKED` — the instance the server reports `affectedRows = 1` to owns the
|
||||
* row, every other sweeper gets 0 and moves on. No transaction to hold open and
|
||||
* no MariaDB version floor.
|
||||
*/
|
||||
async function claimStart(id, owner, leaseUntil) {
|
||||
const result = await query(
|
||||
`UPDATE event_runs
|
||||
SET status = 'starting', claimed_by = ?, claim_expires_at = ?,
|
||||
started_at = COALESCE(started_at, NOW())
|
||||
WHERE id = ? AND status = 'scheduled'`,
|
||||
[owner, leaseUntil, id],
|
||||
)
|
||||
return Number(result?.affectedRows || 0) === 1
|
||||
}
|
||||
|
||||
/**
|
||||
* Take a lease on a run already in flight, so one tick works on it at a time.
|
||||
*
|
||||
* Unlike `claimStart` this does not change `status` — the run is already
|
||||
* `starting`, `running` or `ending`, and what is being claimed is the right to
|
||||
* advance it.
|
||||
*
|
||||
* **A live lease is not re-enterable, not even by the process that took it**, and
|
||||
* that is the whole point rather than an oversight. `setInterval` fires the next
|
||||
* tick whether or not the last one has returned, so an owner-matches escape
|
||||
* clause here would let one process advance one run twice at once — which is
|
||||
* precisely the overrun the plan says the CAS is meant to protect against. A run
|
||||
* this process still holds is a run this process is still working on; the tick
|
||||
* skips it, and `releaseClaim` below is what ends that in the ordinary case.
|
||||
*/
|
||||
async function claimTick(id, owner, leaseUntil, now) {
|
||||
const result = await query(
|
||||
`UPDATE event_runs
|
||||
SET claimed_by = ?, claim_expires_at = ?
|
||||
WHERE id = ?
|
||||
AND status IN ('starting','running','ending')
|
||||
AND (claim_expires_at IS NULL OR claim_expires_at < ?)`,
|
||||
[owner, leaseUntil, id, now],
|
||||
)
|
||||
return Number(result?.affectedRows || 0) === 1
|
||||
}
|
||||
|
||||
/**
|
||||
* Give a still-in-flight run back, so the next tick can pick it up at once.
|
||||
*
|
||||
* A run left parked on a GM cue, or waiting out a `core.wait`, is not finished
|
||||
* and must not carry a lease: without this the run would be unadvanceable until
|
||||
* the lease expired, which would turn every wait into `max(wait, leaseMs)`.
|
||||
* Scoped to `claimed_by = ?` so a process can only release its own claim.
|
||||
*/
|
||||
async function releaseClaim(id, owner) {
|
||||
const result = await query(
|
||||
'UPDATE event_runs SET claimed_by = NULL, claim_expires_at = NULL WHERE id = ? AND claimed_by = ?',
|
||||
[id, owner],
|
||||
)
|
||||
return Number(result?.affectedRows || 0) === 1
|
||||
}
|
||||
|
||||
/**
|
||||
* A guarded status transition: `from -> to`, and only from `from`.
|
||||
*
|
||||
* Every move the runner makes goes through here rather than through a bare
|
||||
* UPDATE, so "did this transition actually happen" is answerable at each call
|
||||
* site. A `false` is not an error — it is another worker, or this run having been
|
||||
* cancelled from the admin surface between the read and the write, which is a
|
||||
* race Phase 3's live controls make ordinary.
|
||||
*/
|
||||
async function transition(id, from, to, { phase, error, clearClaim = false } = {}) {
|
||||
const sets = ['status = ?']
|
||||
const args = [to]
|
||||
if (phase !== undefined) {
|
||||
sets.push('current_phase = ?')
|
||||
args.push(phase)
|
||||
}
|
||||
if (error !== undefined) {
|
||||
sets.push('last_error = ?')
|
||||
args.push(error === null ? null : String(error).slice(0, 500))
|
||||
}
|
||||
if (TERMINAL.includes(to)) sets.push('ended_at = COALESCE(ended_at, NOW())')
|
||||
if (clearClaim || TERMINAL.includes(to)) sets.push('claimed_by = NULL', 'claim_expires_at = NULL')
|
||||
|
||||
const froms = Array.isArray(from) ? from : [from]
|
||||
const result = await query(
|
||||
`UPDATE event_runs SET ${sets.join(', ')}
|
||||
WHERE id = ? AND status IN (${froms.map(() => '?').join(',')})`,
|
||||
[...args, id, ...froms],
|
||||
)
|
||||
return Number(result?.affectedRows || 0) === 1
|
||||
}
|
||||
|
||||
/**
|
||||
* Set health without touching status (§E).
|
||||
*
|
||||
* The two columns are separate because a run can be genuinely running and
|
||||
* degraded at once — announcements landing, world writes parked — and one column
|
||||
* cannot say both. Guarded on the current value so a tick that re-observes the
|
||||
* same degradation does not restamp `updated_at`.
|
||||
*/
|
||||
async function setHealth(id, health) {
|
||||
const result = await query('UPDATE event_runs SET health = ? WHERE id = ? AND health <> ?', [
|
||||
health,
|
||||
id,
|
||||
health,
|
||||
])
|
||||
return Number(result?.affectedRows || 0) === 1
|
||||
}
|
||||
|
||||
/**
|
||||
* Runs whose start instant passed more than their own grace window ago (§E, §L).
|
||||
*
|
||||
* The window is per definition, so the comparison is against `grace_seconds` on
|
||||
* the joined row rather than against a constant here: an event whose announcement
|
||||
* gives a fifteen-minute window and one that must start on the second are the
|
||||
* same query with different data.
|
||||
*
|
||||
* Only `scheduled` runs qualify. A run that reached `starting` has begun, and
|
||||
* "began and then stalled" is a different fact from "never began" — conflating
|
||||
* them would let `missed` describe a run that had already announced itself.
|
||||
*/
|
||||
const findMissed = async (now, limit = 100) => {
|
||||
const n = Math.min(Math.max(Number(limit) || 100, 1), 500)
|
||||
return (
|
||||
await query(
|
||||
`SELECT r.* FROM event_runs r
|
||||
JOIN event_definitions d ON d.id = r.definition_id
|
||||
WHERE r.status = 'scheduled'
|
||||
AND r.scheduled_for + INTERVAL d.grace_seconds SECOND < ?
|
||||
ORDER BY r.scheduled_for
|
||||
LIMIT ${n}`,
|
||||
[now],
|
||||
)
|
||||
).map(hydrate)
|
||||
}
|
||||
|
||||
/**
|
||||
* Is another run holding this concurrency key?
|
||||
*
|
||||
* The org lead's answer for a held key (2026-09-02) is to leave the run
|
||||
* `scheduled` and let the grace window decide, so this is a READ rather than a
|
||||
* claim: the caller holds off, logs which run holds the key, and tries again next
|
||||
* tick. `idx_evrun_concurrency (concurrency_key, status)` is this query, and a
|
||||
* NULL key is skipped by that index — which is right, because a definition with
|
||||
* no key never contends.
|
||||
*/
|
||||
const concurrencyHolder = async (key, exceptRunId) => {
|
||||
if (!key) return null
|
||||
const [row] = await query(
|
||||
`SELECT id, status, definition_id FROM event_runs
|
||||
WHERE concurrency_key = ?
|
||||
AND id <> ?
|
||||
AND status IN ('starting','running','paused','ending')
|
||||
ORDER BY id LIMIT 1`,
|
||||
[key, exceptRunId],
|
||||
)
|
||||
return row || null
|
||||
}
|
||||
|
||||
/**
|
||||
* Recover runs whose claim outlived the process that took it.
|
||||
*
|
||||
* **It does not change status and it touches no counter.** All it releases is the
|
||||
* lease; the run stays exactly where it was and the next tick picks it up through
|
||||
* `findDue`. This is Engagement Phase 14's lesson applied one table over: a sweep
|
||||
* that returned a stale row to its start state made the attempt ceiling
|
||||
* unreachable, so the row cycled forever, never terminal, and therefore never
|
||||
* eligible for any retention sweep.
|
||||
*/
|
||||
const reclaimStale = async (now) => {
|
||||
const result = await query(
|
||||
`UPDATE event_runs SET claimed_by = NULL, claim_expires_at = NULL
|
||||
WHERE status IN ('starting','running','ending')
|
||||
AND claim_expires_at IS NOT NULL
|
||||
AND claim_expires_at < ?`,
|
||||
[now],
|
||||
)
|
||||
return Number(result?.affectedRows || 0)
|
||||
}
|
||||
|
||||
/** Terminal runs that ended before `before` — what the log retention sweep walks. */
|
||||
const terminalBefore = async (before, limit = 500) => {
|
||||
const n = Math.min(Math.max(Number(limit) || 500, 1), 5000)
|
||||
return (
|
||||
await query(
|
||||
`SELECT id FROM event_runs
|
||||
WHERE status IN (${TERMINAL.map(() => '?').join(',')})
|
||||
AND COALESCE(ended_at, updated_at) < ?
|
||||
ORDER BY id LIMIT ${n}`,
|
||||
[...TERMINAL, before],
|
||||
)
|
||||
).map((r) => Number(r.id))
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
list,
|
||||
getById,
|
||||
materialise,
|
||||
findOccurrence,
|
||||
countActiveForDefinition,
|
||||
findDue,
|
||||
findMissed,
|
||||
claimStart,
|
||||
claimTick,
|
||||
releaseClaim,
|
||||
transition,
|
||||
setHealth,
|
||||
concurrencyHolder,
|
||||
reclaimStale,
|
||||
terminalBefore,
|
||||
TERMINAL,
|
||||
}
|
||||
|
||||
@@ -15,6 +15,7 @@ const engagementRetentionPrune = require('./utils/engagementRetentionPrune')
|
||||
const teamForumUploadSweep = require('./utils/teamForumUploadSweep')
|
||||
const teamDigestWorker = require('./utils/teamDigestWorker')
|
||||
const engagementWorker = require('./utils/engagementWorker')
|
||||
const eventRunner = require('./utils/eventRunner')
|
||||
const { ensureSchema, close } = require('./utils/db')
|
||||
const { seedDefaults, createInitialAdminFromEnv } = require('../db/seed')
|
||||
const settings = require('./model/settings/settings.model')
|
||||
@@ -169,6 +170,11 @@ async function start() {
|
||||
// enables a rule: core seeds none and `enabled` defaults to 0.
|
||||
engagementWorker.start()
|
||||
|
||||
// Advance scheduled events (EVENTS.md §E). Materialise, advance, drain, and the
|
||||
// run-log retention sweep. No-op until an admin publishes a definition and
|
||||
// starts a run: core ships no event definitions.
|
||||
eventRunner.start()
|
||||
|
||||
setupShutdown(server, internalServer)
|
||||
}
|
||||
|
||||
@@ -192,6 +198,7 @@ function setupShutdown(server, internalServer) {
|
||||
teamForumUploadSweep.stop() // stop the forum upload sweep
|
||||
teamDigestWorker.stop() // stop the Team forum digest timer
|
||||
engagementWorker.stop() // stop the engagement outbox worker
|
||||
eventRunner.stop() // stop the event runner
|
||||
server.close(() => log.info('http server closed'))
|
||||
if (internalServer) internalServer.close(() => log.info('internal http server closed'))
|
||||
try {
|
||||
|
||||
510
server/src/utils/eventRunner.js
Normal file
510
server/src/utils/eventRunner.js
Normal file
@@ -0,0 +1,510 @@
|
||||
// ── The event runner ───────────────────────────────────────────────────────
|
||||
//
|
||||
// EVENTS.md §E, and Phase 2 of EVENTS_PLAN.md. The eighth poller: same
|
||||
// `setInterval` + `unref()` + `stop()` shape as `announceWorker`, the three Team
|
||||
// sweepers and `engagementWorker`, wired into `server.js` beside them. In the
|
||||
// **website** process rather than the bot, which cannot load module code.
|
||||
//
|
||||
// Its tick does four things, in this order:
|
||||
//
|
||||
// 1. **reclaim** — release leases whose holder died, never touching `attempts`
|
||||
// 2. **materialise** — sweep occurrences past their grace window into `missed`
|
||||
// 3. **advance** — claim each due run and move it through its phases
|
||||
// 4. **prune** — the `event_run_log` retention sweep, on its own long clock
|
||||
//
|
||||
// **What "materialise" means in this phase.** §E's tick materialises due
|
||||
// occurrences from a recurrence; the spec validator accepts `kind: 'manual'`
|
||||
// alone until Phase 4, so there is no recurrence to expand and the only
|
||||
// occurrences that exist are the ones an admin created. What that leaves for this
|
||||
// leg is the half that is already real and already needed: the grace window. A
|
||||
// run whose instant passed while the process was down does not start late and
|
||||
// silently — it becomes `missed`, which is terminal and which a human can see
|
||||
// (§L). Phase 4 adds the expansion above it.
|
||||
//
|
||||
// **Two properties this file must not lose**, both already paid for once on this
|
||||
// codebase:
|
||||
//
|
||||
// - **A reclaim never resets `attempts`.** Engagement Phase 14's defect: a sweep
|
||||
// that returned every stale row to its start state made the attempt ceiling
|
||||
// unreachable, so the row cycled forever, never terminal, and therefore never
|
||||
// eligible for any retention sweep.
|
||||
// - **The unique index, not the claim, is what prevents a double run.** The claim
|
||||
// decides *who* advances an occurrence; `uq_evrun_occurrence` is what stops two
|
||||
// of them existing. Neither substitutes for the other.
|
||||
//
|
||||
// §N4 settled this deployment as single-instance, so the `--scale app=2` rig the
|
||||
// engagement workstream used is deliberately not built. Every claim path is here
|
||||
// exactly as §E specifies anyway, because the multi-instance case is not what
|
||||
// they are for on a single container: the CAS is what protects a tick that
|
||||
// overran into the next one, and the lease and its reclaim are what recover a
|
||||
// step whose process died mid-dispatch. Both happen with one app.
|
||||
|
||||
const os = require('os')
|
||||
|
||||
const runsDb = require('../model/events/eventRuns.db')
|
||||
const stepsDb = require('../model/events/eventRunSteps.db')
|
||||
const logDb = require('../model/events/eventRunLog.db')
|
||||
const versionsDb = require('../model/events/eventVersions.db')
|
||||
const registries = require('../modules/registries')
|
||||
const { dispatchStep } = require('../events/dispatch')
|
||||
const log = require('./logger')('event-runner')
|
||||
|
||||
const POLL_MS = Number(process.env.EVENT_POLL_MS) || 15_000
|
||||
|
||||
// How many runs one sweep looks at, and how many steps it will drain from one
|
||||
// run. Bounds rather than targets: the tick runs again in POLL_MS, and an
|
||||
// unbounded batch is how a backlog turns one tick into a stall. The step bound
|
||||
// also caps how long one run can hold the tick, which is what keeps a
|
||||
// forty-step phase from starving every other run on the deployment.
|
||||
const RUN_BATCH = Number(process.env.EVENT_RUN_BATCH) || 50
|
||||
const STEPS_PER_TICK = Number(process.env.EVENT_STEPS_PER_TICK) || 25
|
||||
|
||||
// A step's retries (org lead, 2026-09-02). §L specifies `retry(n) -> skip |
|
||||
// pause | abort_run` and names no `n`; it lives here, as one number an operator
|
||||
// can change, rather than in a column no authoring surface would ever show.
|
||||
//
|
||||
// Flat backoff rather than exponential, for the reason the outbox's is flat:
|
||||
// `due_at` is also the event's own clock, and a doubling backoff pushes a step
|
||||
// arbitrarily far past the moment the event was about.
|
||||
const MAX_ATTEMPTS = Number(process.env.EVENT_STEP_MAX_ATTEMPTS) || 3
|
||||
const RETRY_MS = Number(process.env.EVENT_STEP_RETRY_MS) || 60_000
|
||||
|
||||
// How long a claim may look alive before the reclaim takes it back. The run
|
||||
// lease has to outlast a whole tick's work on one run; a step's is computed from
|
||||
// its own action's `budgetMs` (see `leaseFor`), because a registry that lets an
|
||||
// action declare an hour would otherwise have its steps reclaimed and
|
||||
// re-dispatched fifty-nine minutes before they answered.
|
||||
const RUN_LEASE_MS = Number(process.env.EVENT_RUN_LEASE_MS) || 15 * 60 * 1000
|
||||
const STEP_LEASE_MARGIN_MS = 60_000
|
||||
|
||||
// The log retention horizon, and how often the sweep runs. It is folded into
|
||||
// this tick rather than given a ninth timer because it shares the tick's only
|
||||
// real dependency — a database — and a sweep that runs four times a day does not
|
||||
// need an interval of its own. Only TERMINAL runs are ever eligible, which is the
|
||||
// rule Engagement Phase 14 arrived at.
|
||||
const LOG_RETENTION_DAYS = Number(process.env.EVENT_LOG_RETENTION_DAYS) || 90
|
||||
const PRUNE_EVERY_MS = 6 * 60 * 60 * 1000
|
||||
|
||||
// Who this process is, for `claimed_by`. Host and pid, so a stranded claim in the
|
||||
// table names the thing that stranded it.
|
||||
const OWNER = `${os.hostname()}:${process.pid}`.slice(0, 64)
|
||||
|
||||
const RUN_TERMINAL = runsDb.TERMINAL
|
||||
|
||||
/** The lease a step's dispatch gets: its action's own budget, plus a margin. */
|
||||
function leaseFor(step, now) {
|
||||
const action = registries.eventAction(step.action_id)
|
||||
const budget = action?.budgetMs || 10_000
|
||||
return new Date(now.getTime() + budget + STEP_LEASE_MARGIN_MS)
|
||||
}
|
||||
|
||||
// ── The disposition of a step that has run out of road ─────────────────────
|
||||
//
|
||||
// **All three dispositions write the step `failed`.** `on_failure` says what
|
||||
// happens to the RUN, not what happened to the step, and a step that was
|
||||
// attempted three times and never worked is `failed` under every one of them.
|
||||
// `skipped` is reserved for a step a human skipped from the run console (Phase
|
||||
// 3) — a status that meant both "nobody ran this" and "this failed and we moved
|
||||
// on" would make the run console's summary line unreadable.
|
||||
async function applyFailure(run, step, error) {
|
||||
await stepsDb.finish(step.id, 'failed', error)
|
||||
await logDb.write({
|
||||
runId: run.id,
|
||||
stepId: step.id,
|
||||
kind: 'step.status',
|
||||
phase: step.phase,
|
||||
detail: { to: 'failed', action: step.action_id, attempts: step.attempts + 1, onFailure: step.on_failure, error },
|
||||
})
|
||||
|
||||
// A run that lost a step is degraded whatever happens next. Health is not
|
||||
// status (§E): a run can be genuinely running and degraded at once, and the
|
||||
// admin surface needs to say so without claiming the run stopped.
|
||||
if (await runsDb.setHealth(run.id, 'degraded')) {
|
||||
await logDb.write({ runId: run.id, kind: 'run.health', detail: { to: 'degraded', because: step.action_id } })
|
||||
}
|
||||
|
||||
if (step.on_failure === 'abort_run') {
|
||||
// §L: the disposition for `irreversible`. Nothing further is dispatched, and
|
||||
// the pending steps are cancelled rather than left looking due forever.
|
||||
const cancelled = await stepsDb.cancelPending(run.id)
|
||||
await runsDb.transition(run.id, ['starting', 'running', 'ending'], 'failed', { error })
|
||||
await logDb.write({
|
||||
runId: run.id,
|
||||
kind: 'run.status',
|
||||
phase: step.phase,
|
||||
detail: { to: 'failed', because: step.action_id, cancelledSteps: cancelled },
|
||||
})
|
||||
return 'stop'
|
||||
}
|
||||
|
||||
if (step.on_failure === 'pause') {
|
||||
// The default for `change`, and the right one when the world is half-altered:
|
||||
// stop advancing and wait for a human. `paused` is excluded from `findDue`,
|
||||
// so nothing here picks it up again — Phase 3's resume control is the only
|
||||
// thing that moves it.
|
||||
await runsDb.transition(run.id, ['starting', 'running'], 'paused', { error })
|
||||
await logDb.write({
|
||||
runId: run.id,
|
||||
kind: 'run.status',
|
||||
phase: step.phase,
|
||||
detail: { to: 'paused', because: step.action_id },
|
||||
})
|
||||
return 'stop'
|
||||
}
|
||||
|
||||
// 'skip': the run carries on, degraded, and the failure is on the record.
|
||||
return 'continue'
|
||||
}
|
||||
|
||||
/**
|
||||
* Claim one step, dispatch it, and record what came back.
|
||||
*
|
||||
* Answers `'continue'` (the phase may proceed), `'stop'` (it may not, for now or
|
||||
* ever) or `'taken'` (somebody else claimed it first).
|
||||
*/
|
||||
async function drainStep(run, step, now, carry = {}) {
|
||||
if (!(await stepsDb.claim(step.id, OWNER, leaseFor(step, now), now))) return 'taken'
|
||||
|
||||
const result = await dispatchStep(step, { run })
|
||||
|
||||
if (result.actionVersionDrift) {
|
||||
await logDb.write({
|
||||
runId: run.id,
|
||||
stepId: step.id,
|
||||
kind: 'step.status',
|
||||
phase: step.phase,
|
||||
detail: { action: step.action_id, versionDrift: result.actionVersionDrift },
|
||||
})
|
||||
}
|
||||
|
||||
if (result.outcome === 'parked') {
|
||||
// The GM cue. The step stays `running` with a NULL lease: genuinely in
|
||||
// flight, nothing holding it, so the stale reclaim passes it by and a cue
|
||||
// posted on Friday is still waiting on Monday. Phase 3's confirm control is
|
||||
// what ends it.
|
||||
await stepsDb.park(step.id, null)
|
||||
await logDb.write({
|
||||
runId: run.id,
|
||||
stepId: step.id,
|
||||
kind: 'step.parked',
|
||||
phase: step.phase,
|
||||
detail: { action: step.action_id, params: step.params },
|
||||
})
|
||||
return 'stop'
|
||||
}
|
||||
|
||||
if (result.outcome === 'done') {
|
||||
await stepsDb.finish(step.id, 'done', null)
|
||||
if (result.holdSeconds > 0) {
|
||||
// `core.wait`, and any module action that answers `holdFor`. The pause is
|
||||
// the NEXT step's `due_at` and it is set here, by the runner, because a
|
||||
// `perform()` that slept would hold its claim for the duration.
|
||||
const until = new Date(now.getTime() + result.holdSeconds * 1000)
|
||||
if (!(await stepsDb.holdNext(run.id, step.phase, step.seq, until))) {
|
||||
// **Nothing after it in this phase**, which is the case a wait written as
|
||||
// the last step of a phase produces. Dropping the hold here would make
|
||||
// "announce, wait five minutes, then the next phase" start the next phase
|
||||
// at once — a wait that silently meant nothing. The later phase's steps do
|
||||
// not exist yet, so the instant is carried out to `advanceRun` and applied
|
||||
// when they are materialised.
|
||||
carry.holdUntil = until
|
||||
}
|
||||
}
|
||||
await logDb.write({
|
||||
runId: run.id,
|
||||
stepId: step.id,
|
||||
kind: 'step.status',
|
||||
phase: step.phase,
|
||||
detail: { to: 'done', action: step.action_id, holdSeconds: result.holdSeconds || 0 },
|
||||
})
|
||||
return 'continue'
|
||||
}
|
||||
|
||||
if (result.outcome === 'retry' && step.attempts + 1 < MAX_ATTEMPTS) {
|
||||
await stepsDb.reschedule(step.id, new Date(now.getTime() + RETRY_MS), result.error)
|
||||
await logDb.write({
|
||||
runId: run.id,
|
||||
stepId: step.id,
|
||||
kind: 'step.retry',
|
||||
phase: step.phase,
|
||||
detail: { action: step.action_id, attempt: step.attempts + 1, of: MAX_ATTEMPTS, error: result.error },
|
||||
})
|
||||
// Degraded from the FIRST retry, not from the eventual failure. An event
|
||||
// whose announcements are landing on the second attempt is having trouble
|
||||
// now, and that is when an operator wants to know.
|
||||
if (await runsDb.setHealth(run.id, 'degraded')) {
|
||||
await logDb.write({ runId: run.id, kind: 'run.health', detail: { to: 'degraded', because: step.action_id } })
|
||||
}
|
||||
return 'stop'
|
||||
}
|
||||
|
||||
// Terminal, or transient with the attempts spent. Same disposition either way:
|
||||
// §L's `retry(n) -> ...` has arrived at the arrow.
|
||||
return applyFailure(run, step, result.error)
|
||||
}
|
||||
|
||||
/**
|
||||
* Advance one claimed run as far as it will go this tick.
|
||||
*
|
||||
* The loop is bounded by `STEPS_PER_TICK` and exits on the first thing it cannot
|
||||
* get past — a parked step, a step whose `due_at` is in the future, a step
|
||||
* somebody else holds, or a phase that is not finished.
|
||||
*/
|
||||
async function advanceRun(run, now) {
|
||||
const version = await versionsDb.getById(run.version_id)
|
||||
const phases = version?.spec?.phases
|
||||
if (!Array.isArray(phases) || !phases.length) {
|
||||
// The pinned version is unreadable. `version_id`'s foreign key RESTRICTs
|
||||
// precisely so this cannot be a deleted row, so it is corruption rather than
|
||||
// an ordinary race — terminal, named, and not retried.
|
||||
await runsDb.transition(run.id, ['scheduled', 'starting', 'running', 'ending'], 'failed', {
|
||||
error: 'the pinned version has no phases',
|
||||
})
|
||||
await logDb.write({ runId: run.id, kind: 'run.status', detail: { to: 'failed', because: 'pinned version has no phases' } })
|
||||
return 'failed'
|
||||
}
|
||||
|
||||
let phaseKey = run.current_phase
|
||||
|
||||
if (run.status === 'starting') {
|
||||
// Entering the first phase. `materialisePhase` is INSERT IGNORE against
|
||||
// `uq_evstep_slot`, so doing it again over the rows Phase 1's `create()`
|
||||
// already wrote is a no-op — which is what makes recovery from a process that
|
||||
// died between the claim and here uneventful.
|
||||
const first = phases[0]
|
||||
await stepsDb.materialisePhase(run.id, first.key, first.steps || [])
|
||||
if (!(await runsDb.transition(run.id, 'starting', 'running', { phase: first.key }))) return 'taken'
|
||||
phaseKey = first.key
|
||||
await logDb.write({ runId: run.id, kind: 'run.status', phase: first.key, detail: { from: 'starting', to: 'running' } })
|
||||
}
|
||||
|
||||
if (run.status === 'ending') {
|
||||
// A run that reached the wind-down and then lost its process. Phase 8 puts
|
||||
// cleanup here; until then `ending` is a state a run passes through rather
|
||||
// than one it does work in, and completing it is the whole recovery.
|
||||
await runsDb.transition(run.id, 'ending', 'completed')
|
||||
await logDb.write({ runId: run.id, kind: 'run.status', detail: { from: 'ending', to: 'completed' } })
|
||||
return 'completed'
|
||||
}
|
||||
|
||||
// A hold a `core.wait` could not place because nothing followed it in its own
|
||||
// phase. It crosses the phase boundary with the run rather than being dropped.
|
||||
const carry = {}
|
||||
|
||||
for (let n = 0; n < STEPS_PER_TICK; n++) {
|
||||
const phaseIndex = phases.findIndex((p) => p.key === phaseKey)
|
||||
if (phaseIndex < 0) {
|
||||
await runsDb.transition(run.id, ['running'], 'failed', { error: `phase "${phaseKey}" is not in the pinned version` })
|
||||
await logDb.write({ runId: run.id, kind: 'run.status', detail: { to: 'failed', because: `unknown phase "${phaseKey}"` } })
|
||||
return 'failed'
|
||||
}
|
||||
|
||||
const step = await stepsDb.nextOpenStep(run.id, phaseKey)
|
||||
|
||||
if (step && step.status === 'running') return 'in-flight' // parked, or somebody's dispatch
|
||||
if (step && step.due_at && new Date(step.due_at) > now) return 'waiting' // behind a core.wait
|
||||
|
||||
if (step) {
|
||||
const outcome = await drainStep({ ...run, current_phase: phaseKey }, step, now, carry)
|
||||
if (outcome === 'continue') continue
|
||||
return outcome === 'taken' ? 'taken' : 'stopped'
|
||||
}
|
||||
|
||||
// Every step of this phase is terminal.
|
||||
await logDb.write({ runId: run.id, kind: 'phase.completed', phase: phaseKey, detail: { index: phaseIndex } })
|
||||
|
||||
const next = phases[phaseIndex + 1]
|
||||
if (!next) {
|
||||
// §E's `ending` exists for the reason `sending` does in the outbox — it is
|
||||
// what a claim sets — so the run passes through it even though Phase 2 has
|
||||
// no cleanup to do there. Phase 8 is what gives it work.
|
||||
if (!(await runsDb.transition(run.id, 'running', 'ending'))) return 'taken'
|
||||
await logDb.write({ runId: run.id, kind: 'run.status', phase: phaseKey, detail: { from: 'running', to: 'ending' } })
|
||||
await runsDb.transition(run.id, 'ending', 'completed')
|
||||
await logDb.write({ runId: run.id, kind: 'run.status', detail: { from: 'ending', to: 'completed' } })
|
||||
return 'completed'
|
||||
}
|
||||
|
||||
await stepsDb.materialisePhase(run.id, next.key, next.steps || [])
|
||||
if (carry.holdUntil) {
|
||||
// `seq > -1` is the first step of the phase just created. Applied after
|
||||
// materialisation because that is the first moment there is a row to hold.
|
||||
await stepsDb.holdNext(run.id, next.key, -1, carry.holdUntil)
|
||||
carry.holdUntil = null
|
||||
}
|
||||
// `running -> running` is not a no-op: it is a guarded write of
|
||||
// `current_phase` that fails if the run stopped being `running` underneath
|
||||
// this tick, which is what a cancel from the admin surface looks like.
|
||||
if (!(await runsDb.transition(run.id, 'running', 'running', { phase: next.key }))) return 'taken'
|
||||
phaseKey = next.key
|
||||
await logDb.write({ runId: run.id, kind: 'phase.entered', phase: next.key, detail: { steps: (next.steps || []).length } })
|
||||
}
|
||||
|
||||
return 'bounded' // more to do; the next tick picks it up
|
||||
}
|
||||
|
||||
/** Claim one due run, work it, and hand the lease back if it is still in flight. */
|
||||
async function processRun(run, now = new Date()) {
|
||||
if (run.status === 'scheduled') {
|
||||
const holder = await runsDb.concurrencyHolder(run.concurrency_key, run.id)
|
||||
if (holder) {
|
||||
// The org lead's answer for a held key (2026-09-02): hold at `scheduled`
|
||||
// and let the grace window decide. Nothing is destroyed, nothing starts
|
||||
// silently late, and if the holder outlasts the window the missed sweep
|
||||
// makes this run terminal and visible.
|
||||
//
|
||||
// Logged only when the reason CHANGES. A blocked run is re-examined every
|
||||
// tick, and a line per tick for the length of a grace window would bury the
|
||||
// one line that matters under a thousand identical ones.
|
||||
const message = `held: run ${holder.id} has concurrency key "${run.concurrency_key}"`
|
||||
if (run.last_error !== message) {
|
||||
await runsDb.transition(run.id, 'scheduled', 'scheduled', { error: message })
|
||||
await logDb.write({
|
||||
runId: run.id,
|
||||
kind: 'run.blocked',
|
||||
detail: { concurrencyKey: run.concurrency_key, heldBy: holder.id, holderStatus: holder.status },
|
||||
})
|
||||
}
|
||||
return 'blocked'
|
||||
}
|
||||
|
||||
if (!(await runsDb.claimStart(run.id, OWNER, new Date(now.getTime() + RUN_LEASE_MS)))) return 'taken'
|
||||
await logDb.write({ runId: run.id, kind: 'run.status', detail: { from: 'scheduled', to: 'starting' } })
|
||||
run = { ...run, status: 'starting' }
|
||||
} else if (!(await runsDb.claimTick(run.id, OWNER, new Date(now.getTime() + RUN_LEASE_MS), now))) {
|
||||
return 'taken'
|
||||
}
|
||||
|
||||
try {
|
||||
const outcome = await advanceRun(run, now)
|
||||
// A run that is not finished must not keep its lease: it would be
|
||||
// unadvanceable until the lease expired, which would turn every `core.wait`
|
||||
// into `max(wait, RUN_LEASE_MS)`. A terminal transition already cleared it.
|
||||
if (!['completed', 'failed'].includes(outcome)) await runsDb.releaseClaim(run.id, OWNER)
|
||||
return outcome
|
||||
} catch (err) {
|
||||
await runsDb.releaseClaim(run.id, OWNER)
|
||||
throw err
|
||||
}
|
||||
}
|
||||
|
||||
/** Occurrences that passed their own grace window while nothing was running (§L). */
|
||||
async function sweepMissed(now) {
|
||||
const missed = await runsDb.findMissed(now)
|
||||
let n = 0
|
||||
for (const run of missed) {
|
||||
if (await runsDb.transition(run.id, 'scheduled', 'missed', { error: 'the grace window passed' })) {
|
||||
await stepsDb.cancelPending(run.id)
|
||||
await logDb.write({
|
||||
runId: run.id,
|
||||
kind: 'run.status',
|
||||
detail: { to: 'missed', scheduledFor: run.scheduled_for },
|
||||
})
|
||||
n += 1
|
||||
}
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
let lastPruneAt = 0
|
||||
|
||||
async function prune(now) {
|
||||
if (now.getTime() - lastPruneAt < PRUNE_EVERY_MS) return 0
|
||||
lastPruneAt = now.getTime()
|
||||
const before = new Date(now.getTime() - LOG_RETENTION_DAYS * 24 * 60 * 60 * 1000)
|
||||
const deleted = await logDb.pruneTerminal(before)
|
||||
if (deleted) log.info('event run log pruned', { deleted, before, retentionDays: LOG_RETENTION_DAYS })
|
||||
return deleted
|
||||
}
|
||||
|
||||
async function tick(now = new Date()) {
|
||||
try {
|
||||
await runsDb.reclaimStale(now)
|
||||
await stepsDb.reclaimStale(now, MAX_ATTEMPTS)
|
||||
} catch (err) {
|
||||
log.error('failed to reclaim stale claims', { message: err.message })
|
||||
}
|
||||
|
||||
try {
|
||||
await sweepMissed(now)
|
||||
} catch (err) {
|
||||
log.error('missed sweep failed', { message: err.message })
|
||||
}
|
||||
|
||||
let due
|
||||
try {
|
||||
due = await runsDb.findDue(now, RUN_BATCH)
|
||||
} catch (err) {
|
||||
log.error('failed to load due runs', { message: err.message })
|
||||
return
|
||||
}
|
||||
|
||||
const counts = {}
|
||||
for (const run of due || []) {
|
||||
try {
|
||||
const outcome = await processRun(run, now)
|
||||
counts[outcome] = (counts[outcome] || 0) + 1
|
||||
} catch (err) {
|
||||
log.error('run failed', { run: run.id, message: err.message })
|
||||
}
|
||||
}
|
||||
if (due && due.length) log.info('event runs swept', { due: due.length, ...counts })
|
||||
|
||||
try {
|
||||
await prune(now)
|
||||
} catch (err) {
|
||||
log.error('log prune failed', { message: err.message })
|
||||
}
|
||||
}
|
||||
|
||||
let timer = null
|
||||
// Guards against this process running two ticks over the same runs at once.
|
||||
// `setInterval` fires whether or not the last callback returned, and the CAS
|
||||
// alone does not cover it now that a live lease is not re-enterable by its own
|
||||
// owner — a second tick would simply find every run claimed and do nothing
|
||||
// useful, one query at a time, for as long as the first one ran.
|
||||
let ticking = false
|
||||
|
||||
function start() {
|
||||
if (timer) return timer
|
||||
timer = setInterval(() => {
|
||||
if (ticking) {
|
||||
log.warn('event tick still running; skipping this interval')
|
||||
return
|
||||
}
|
||||
ticking = true
|
||||
tick()
|
||||
.catch((err) => log.error('event tick failed', { message: err.message }))
|
||||
.finally(() => {
|
||||
ticking = false
|
||||
})
|
||||
}, POLL_MS)
|
||||
if (timer.unref) timer.unref() // don't keep the event loop alive (tests, shutdown)
|
||||
log.info('event runner started', { pollMs: POLL_MS, owner: OWNER, maxAttempts: MAX_ATTEMPTS })
|
||||
return timer
|
||||
}
|
||||
|
||||
function stop() {
|
||||
if (timer) {
|
||||
clearInterval(timer)
|
||||
timer = null
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
start,
|
||||
stop,
|
||||
tick,
|
||||
processRun,
|
||||
advanceRun,
|
||||
drainStep,
|
||||
sweepMissed,
|
||||
prune,
|
||||
OWNER,
|
||||
POLL_MS,
|
||||
MAX_ATTEMPTS,
|
||||
RETRY_MS,
|
||||
RUN_LEASE_MS,
|
||||
LOG_RETENTION_DAYS,
|
||||
RUN_TERMINAL,
|
||||
}
|
||||
Reference in New Issue
Block a user