feat(events): the minimal admin surface (Phase 3)
Three screens, an Events nav group and the six live run controls Phase 1 left
absent on purpose because nothing was in flight. An admin can now author,
publish, start and watch an event that announces things and cues a human; a
moderator can stop one that is going wrong.
Six controls, not eight. `advance` is absent because a phase today advances when
its steps go terminal — the per-step skip already does that — and Phase 5 is what
gives a phase an advance condition. Cancel takes `{ reason }`, not `{ cleanup }`,
until Phase 8's ledger exists. Every control is a compare-and-set on the status it
may act from, so a console rendered thirty seconds ago cannot act on a run that
has moved.
Fixes a defect in the Phase 2 runner: `advanceRun` drained up to
EVENT_STEPS_PER_TICK steps while only checking the run's status at the top of the
tick, so a pause pressed mid-batch did nothing for up to 24 more steps.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01T6t8mrAWhZU5vnyYgZTMtL
This commit is contained in:
@@ -518,6 +518,15 @@
|
||||
"requireAuth"
|
||||
]
|
||||
},
|
||||
{
|
||||
"method": "POST",
|
||||
"path": "/api/v1/admin/events/runs/:runId/cancel",
|
||||
"handlers": 2,
|
||||
"gates": [
|
||||
"noindex",
|
||||
"requireAuth"
|
||||
]
|
||||
},
|
||||
{
|
||||
"method": "GET",
|
||||
"path": "/api/v1/admin/events/runs/:runId/log",
|
||||
@@ -527,6 +536,51 @@
|
||||
"requireAuth"
|
||||
]
|
||||
},
|
||||
{
|
||||
"method": "POST",
|
||||
"path": "/api/v1/admin/events/runs/:runId/pause",
|
||||
"handlers": 2,
|
||||
"gates": [
|
||||
"noindex",
|
||||
"requireAuth"
|
||||
]
|
||||
},
|
||||
{
|
||||
"method": "POST",
|
||||
"path": "/api/v1/admin/events/runs/:runId/resume",
|
||||
"handlers": 2,
|
||||
"gates": [
|
||||
"noindex",
|
||||
"requireAuth"
|
||||
]
|
||||
},
|
||||
{
|
||||
"method": "POST",
|
||||
"path": "/api/v1/admin/events/runs/:runId/steps/:stepId/confirm",
|
||||
"handlers": 2,
|
||||
"gates": [
|
||||
"noindex",
|
||||
"requireAuth"
|
||||
]
|
||||
},
|
||||
{
|
||||
"method": "POST",
|
||||
"path": "/api/v1/admin/events/runs/:runId/steps/:stepId/retry",
|
||||
"handlers": 2,
|
||||
"gates": [
|
||||
"noindex",
|
||||
"requireAuth"
|
||||
]
|
||||
},
|
||||
{
|
||||
"method": "POST",
|
||||
"path": "/api/v1/admin/events/runs/:runId/steps/:stepId/skip",
|
||||
"handlers": 2,
|
||||
"gates": [
|
||||
"noindex",
|
||||
"requireAuth"
|
||||
]
|
||||
},
|
||||
{
|
||||
"method": "GET",
|
||||
"path": "/api/v1/admin/events/series",
|
||||
|
||||
@@ -229,10 +229,34 @@
|
||||
"method": "GET",
|
||||
"path": "/api/v1/admin/events/runs/:runId"
|
||||
},
|
||||
{
|
||||
"method": "POST",
|
||||
"path": "/api/v1/admin/events/runs/:runId/cancel"
|
||||
},
|
||||
{
|
||||
"method": "GET",
|
||||
"path": "/api/v1/admin/events/runs/:runId/log"
|
||||
},
|
||||
{
|
||||
"method": "POST",
|
||||
"path": "/api/v1/admin/events/runs/:runId/pause"
|
||||
},
|
||||
{
|
||||
"method": "POST",
|
||||
"path": "/api/v1/admin/events/runs/:runId/resume"
|
||||
},
|
||||
{
|
||||
"method": "POST",
|
||||
"path": "/api/v1/admin/events/runs/:runId/steps/:stepId/confirm"
|
||||
},
|
||||
{
|
||||
"method": "POST",
|
||||
"path": "/api/v1/admin/events/runs/:runId/steps/:stepId/retry"
|
||||
},
|
||||
{
|
||||
"method": "POST",
|
||||
"path": "/api/v1/admin/events/runs/:runId/steps/:stepId/skip"
|
||||
},
|
||||
{
|
||||
"method": "GET",
|
||||
"path": "/api/v1/admin/events/series"
|
||||
|
||||
296
server/src/model/events/eventRunControls.model.js
Normal file
296
server/src/model/events/eventRunControls.model.js
Normal file
@@ -0,0 +1,296 @@
|
||||
// ── The live run controls ──────────────────────────────────────────────────
|
||||
//
|
||||
// EVENTS.md §I ("live controls that are honest"), §K and §L. Six of them: pause,
|
||||
// resume and cancel act on a run; confirm, skip and retry act on one step. They
|
||||
// arrive in Phase 3 because Phase 2 is what gave them something to act on — a
|
||||
// run that announces, waits and completes on its own is exactly the run that
|
||||
// needs no control, and a run that paused on a failed world write is the one
|
||||
// that does.
|
||||
//
|
||||
// **Two of §I's six run-level controls are deliberately not here.**
|
||||
// `advance` — force a phase forward — has no honest meaning yet: a phase today
|
||||
// advances when its steps go terminal, and the per-step skip already does that
|
||||
// one step at a time. Phase 5 is what gives a phase an `advance` CONDITION, and
|
||||
// that is the first moment "force it anyway" means something an operator could
|
||||
// predict. `cleanup` needs Phase 8's resource ledger; there is nothing to
|
||||
// revert, so cancel takes `{ reason }` and gains `cleanup` when there is
|
||||
// something for it to do. Both are absent rather than inert, which is the
|
||||
// posture Phase 1 set and Phase 2 kept.
|
||||
//
|
||||
// **Every control is guarded on the status it may act from, and the guard is a
|
||||
// WHERE clause rather than a read-then-write.** A run console rendered thirty
|
||||
// seconds ago describes a run that has since moved — the runner ticks every
|
||||
// fifteen — so a control that checked in JavaScript and then wrote would race
|
||||
// the tick it exists to interrupt. `transition()` and the four step statements
|
||||
// are all compare-and-set, and a `false` from one of them is reported as a 409
|
||||
// naming the status the run is actually in.
|
||||
//
|
||||
// **Who may press them is `admin` + `moderator` (§K, §N2), and it is the widest
|
||||
// gate in this feature on purpose.** Starting a run commits the deployment to
|
||||
// everything the definition contains, unattended — that wants the narrowest gate
|
||||
// there is. Stopping one is incident response at 2am, and it wants the widest.
|
||||
|
||||
const runsDb = require('./eventRuns.db')
|
||||
const stepsDb = require('./eventRunSteps.db')
|
||||
const logDb = require('./eventRunLog.db')
|
||||
|
||||
const MAX_REASON = 500
|
||||
|
||||
const clean = (raw) => {
|
||||
const text = typeof raw === 'string' ? raw.trim() : ''
|
||||
return text ? text.slice(0, MAX_REASON) : null
|
||||
}
|
||||
|
||||
const conflict = (message) => ({ ok: false, status: 409, errors: [message] })
|
||||
|
||||
/** The run, or a 404 shaped the way every other model here shapes one. */
|
||||
async function loadRun(runId) {
|
||||
const run = await runsDb.getById(runId)
|
||||
return run || null
|
||||
}
|
||||
|
||||
/**
|
||||
* A step of THIS run, or null.
|
||||
*
|
||||
* Scoped to the run rather than fetched by id alone: the step id arrives from a
|
||||
* URL under a run id, and a control that acted on a step belonging to a
|
||||
* different run would be a real one — the console's step ids are not secret and
|
||||
* the two paths would otherwise never be compared.
|
||||
*/
|
||||
async function loadStep(runId, stepId) {
|
||||
const step = await stepsDb.getById(stepId)
|
||||
if (!step || Number(step.run_id) !== Number(runId)) return null
|
||||
return step
|
||||
}
|
||||
|
||||
// ── Run-level ─────────────────────────────────────────────────────────────
|
||||
|
||||
/**
|
||||
* Pause a run in flight.
|
||||
*
|
||||
* `starting` and `running` only — §K's "live control of a run **in flight**". A
|
||||
* `scheduled` run has not begun, and the thing to do with an occurrence that
|
||||
* should not happen is cancel it: pausing one would leave a run that is neither
|
||||
* going to start nor visibly abandoned, and resuming it after its grace window
|
||||
* had passed would produce a `missed` from a button labelled resume.
|
||||
*
|
||||
* The claim is cleared with the transition. A tick may be working the run at
|
||||
* this exact moment; it will find its guarded writes returning zero rows and
|
||||
* hand back a lease it no longer holds, both of which are no-ops. What it will
|
||||
* NOT do is dispatch the rest of its batch — `advanceRun` re-reads the status
|
||||
* between steps precisely so this control means what it says.
|
||||
*/
|
||||
async function pause(runId, { reason } = {}, userId = null) {
|
||||
const run = await loadRun(runId)
|
||||
if (!run) return { ok: false, status: 404, errors: ['no such run'] }
|
||||
if (run.status === 'paused') return conflict('this run is already paused')
|
||||
|
||||
const note = clean(reason)
|
||||
if (!(await runsDb.transition(run.id, ['starting', 'running'], 'paused', { clearClaim: true }))) {
|
||||
return conflict(`a ${run.status} run cannot be paused`)
|
||||
}
|
||||
|
||||
await logDb.write({
|
||||
runId: run.id,
|
||||
kind: 'run.status',
|
||||
phase: run.current_phase,
|
||||
detail: { from: run.status, to: 'paused', control: 'pause', by: userId, reason: note },
|
||||
})
|
||||
return { ok: true, run: await runsDb.getById(run.id) }
|
||||
}
|
||||
|
||||
/**
|
||||
* Resume a paused run.
|
||||
*
|
||||
* Where it goes back to is derived rather than remembered: `current_phase` is
|
||||
* set by the transition into `running` and by nothing else, so a paused run that
|
||||
* has one was running and a paused run that has none never got past `starting`.
|
||||
* Both statuses are in `findDue`, so the next tick picks the run up either way,
|
||||
* and there is no fourth column recording what a run was paused *from* — a
|
||||
* column that could disagree with the run's own history.
|
||||
*
|
||||
* **`last_error` is cleared and `health` is not.** The error is what the pause
|
||||
* was about and an operator has just dealt with it; leaving it on the banner
|
||||
* would have a healthy run permanently accused of a failure that is in the log
|
||||
* where it belongs. Health is a different claim — that this run has already had
|
||||
* trouble — and it stays true no matter who pressed resume.
|
||||
*/
|
||||
async function resume(runId, options = {}, userId = null) {
|
||||
const run = await loadRun(runId)
|
||||
if (!run) return { ok: false, status: 404, errors: ['no such run'] }
|
||||
if (run.status !== 'paused') return conflict(`a ${run.status} run is not paused`)
|
||||
|
||||
const to = run.current_phase ? 'running' : 'starting'
|
||||
if (!(await runsDb.transition(run.id, 'paused', to, { error: null }))) {
|
||||
return conflict('this run stopped being paused')
|
||||
}
|
||||
|
||||
await logDb.write({
|
||||
runId: run.id,
|
||||
kind: 'run.status',
|
||||
phase: run.current_phase,
|
||||
detail: { from: 'paused', to, control: 'resume', by: userId },
|
||||
})
|
||||
return { ok: true, run: await runsDb.getById(run.id) }
|
||||
}
|
||||
|
||||
/**
|
||||
* Cancel a run.
|
||||
*
|
||||
* Legal from every non-terminal status including `scheduled`, because "this
|
||||
* event is not happening" is a decision an operator makes before it starts as
|
||||
* often as during it.
|
||||
*
|
||||
* `cancelOpen` then closes out the steps that will never run — the pending ones
|
||||
* and any parked cue. A step with a LIVE lease is left exactly where it is:
|
||||
* something is dispatching it, nothing can recall a command already sent (§L),
|
||||
* and a second writer on that row would race the process that owns it. It
|
||||
* finishes into a cancelled run, which is honest.
|
||||
*/
|
||||
async function cancel(runId, { reason } = {}, userId = null) {
|
||||
const run = await loadRun(runId)
|
||||
if (!run) return { ok: false, status: 404, errors: ['no such run'] }
|
||||
if (runsDb.TERMINAL.includes(run.status)) return conflict(`this run is already ${run.status}`)
|
||||
|
||||
const note = clean(reason)
|
||||
const from = ['scheduled', 'starting', 'running', 'paused', 'ending']
|
||||
if (!(await runsDb.transition(run.id, from, 'cancelled', { error: note || 'cancelled by staff' }))) {
|
||||
return conflict('this run is no longer cancellable')
|
||||
}
|
||||
|
||||
const closed = await stepsDb.cancelOpen(run.id)
|
||||
await logDb.write({
|
||||
runId: run.id,
|
||||
kind: 'run.status',
|
||||
phase: run.current_phase,
|
||||
detail: { from: run.status, to: 'cancelled', control: 'cancel', by: userId, reason: note, cancelledSteps: closed },
|
||||
})
|
||||
return { ok: true, run: await runsDb.getById(run.id), cancelledSteps: closed }
|
||||
}
|
||||
|
||||
// ── Step-level ────────────────────────────────────────────────────────────
|
||||
|
||||
/**
|
||||
* Confirm a parked step — the GM cue's other half.
|
||||
*
|
||||
* `core.cue` posts an instruction and parks: the step stays `running` with a
|
||||
* NULL lease, genuinely in flight with nothing holding it, so no sweep takes it
|
||||
* back and a cue posted on Friday is still waiting on Monday. This is what ends
|
||||
* it, and it is the control that makes the whole system useful before any module
|
||||
* automates anything — a GM does the target-driven part in-client and says so
|
||||
* here.
|
||||
*
|
||||
* The outcome is `done`, not `skipped`: a person saying they did the thing is
|
||||
* the step having succeeded. The note is what they did, and it is kept.
|
||||
*/
|
||||
async function confirmStep(runId, stepId, { note } = {}, userId = null) {
|
||||
const run = await loadRun(runId)
|
||||
if (!run) return { ok: false, status: 404, errors: ['no such run'] }
|
||||
const step = await loadStep(runId, stepId)
|
||||
if (!step) return { ok: false, status: 404, errors: ['no such step on this run'] }
|
||||
|
||||
const text = clean(note)
|
||||
if (!(await stepsDb.confirmParked(step.id, text))) {
|
||||
return conflict(`this step is ${step.status} and is not waiting on anyone`)
|
||||
}
|
||||
|
||||
await logDb.write({
|
||||
runId: run.id,
|
||||
stepId: step.id,
|
||||
kind: 'step.status',
|
||||
phase: step.phase,
|
||||
detail: { to: 'done', action: step.action_id, control: 'confirm', by: userId, note: text },
|
||||
})
|
||||
return { ok: true, step: await stepsDb.getById(step.id) }
|
||||
}
|
||||
|
||||
/**
|
||||
* Skip a step: one that has not started, or a parked cue nobody is going to do.
|
||||
*
|
||||
* This is what the `skipped` status was reserved for (§L) — which is also why
|
||||
* the three `on_failure` dispositions all write `failed` instead. A status
|
||||
* meaning both "a human decided against this" and "this was attempted three
|
||||
* times and never worked" would make the console's summary line unreadable.
|
||||
*
|
||||
* A `failed` step is not skippable and does not need to be: `nextOpenStep`
|
||||
* already passes over one, so resuming a run carries the phase past it.
|
||||
*/
|
||||
async function skipStep(runId, stepId, { reason } = {}, userId = null) {
|
||||
const run = await loadRun(runId)
|
||||
if (!run) return { ok: false, status: 404, errors: ['no such run'] }
|
||||
if (runsDb.TERMINAL.includes(run.status)) return conflict(`this run is ${run.status}`)
|
||||
const step = await loadStep(runId, stepId)
|
||||
if (!step) return { ok: false, status: 404, errors: ['no such step on this run'] }
|
||||
|
||||
const note = clean(reason)
|
||||
if (!(await stepsDb.skipByHuman(step.id, note))) {
|
||||
return conflict(`a ${step.status} step cannot be skipped`)
|
||||
}
|
||||
|
||||
await logDb.write({
|
||||
runId: run.id,
|
||||
stepId: step.id,
|
||||
kind: 'step.status',
|
||||
phase: step.phase,
|
||||
detail: { to: 'skipped', action: step.action_id, control: 'skip', by: userId, reason: note },
|
||||
})
|
||||
return { ok: true, step: await stepsDb.getById(step.id) }
|
||||
}
|
||||
|
||||
/**
|
||||
* Re-queue the failed step a run is stopped at, and resume the run — one action.
|
||||
*
|
||||
* **The two halves are one control because there is no state in which you would
|
||||
* want half of it.** Retry is legal only from `paused`, and a paused run is
|
||||
* paused *at* this step; re-queueing without resuming would leave the run in
|
||||
* precisely the state it was already in, with a button the operator now has to
|
||||
* find. Splitting them would read as honesty and behave as a trap.
|
||||
*
|
||||
* Two guards, and the second is the one worth explaining. The step must be the
|
||||
* furthest one its phase has reached — `lastStartedSeq` — because a `failed`
|
||||
* step under an `on_failure` of `skip` is one the run has already moved PAST.
|
||||
* `nextOpenStep` selects `pending` and `running` only, so the runner steps over
|
||||
* a failed row; re-queueing an earlier one puts a `pending` step behind the
|
||||
* cursor, where it sits for ever.
|
||||
*/
|
||||
async function retryStep(runId, stepId, options = {}, userId = null) {
|
||||
const run = await loadRun(runId)
|
||||
if (!run) return { ok: false, status: 404, errors: ['no such run'] }
|
||||
if (run.status !== 'paused') {
|
||||
return conflict(`a step can only be retried while its run is paused; this run is ${run.status}`)
|
||||
}
|
||||
const step = await loadStep(runId, stepId)
|
||||
if (!step) return { ok: false, status: 404, errors: ['no such step on this run'] }
|
||||
if (step.status !== 'failed') return conflict(`a ${step.status} step cannot be retried`)
|
||||
if (step.phase !== run.current_phase) {
|
||||
return conflict('this step belongs to a phase the run has already left')
|
||||
}
|
||||
|
||||
const furthest = await stepsDb.lastStartedSeq(run.id, step.phase)
|
||||
if (furthest === null || Number(furthest) !== Number(step.seq)) {
|
||||
return conflict('the run is not stopped at this step; only the step a phase is stopped at can be retried')
|
||||
}
|
||||
|
||||
if (!(await stepsDb.requeue(step.id))) return conflict('this step is no longer failed')
|
||||
|
||||
await logDb.write({
|
||||
runId: run.id,
|
||||
stepId: step.id,
|
||||
kind: 'step.status',
|
||||
phase: step.phase,
|
||||
detail: { to: 'pending', action: step.action_id, control: 'retry', by: userId, attemptsReset: step.attempts },
|
||||
})
|
||||
|
||||
const resumed = await resume(runId, {}, userId)
|
||||
return {
|
||||
ok: true,
|
||||
step: await stepsDb.getById(step.id),
|
||||
// A resume that did not take is reported rather than swallowed: the step IS
|
||||
// re-queued either way, and an operator told "retried" about a run that is
|
||||
// still paused would be told something false.
|
||||
resumed: Boolean(resumed.ok),
|
||||
run: resumed.run || (await runsDb.getById(run.id)),
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = { pause, resume, cancel, confirmStep, skipStep, retryStep }
|
||||
@@ -274,6 +274,130 @@ const cancelPending = async (runId) => {
|
||||
return Number(result?.affectedRows || 0)
|
||||
}
|
||||
|
||||
// ── Phase 3: the controls a human works ────────────────────────────────────
|
||||
//
|
||||
// Four statements, and every one of them is guarded on the status it is allowed
|
||||
// to act from rather than trusting the button that was pressed. The run console
|
||||
// decides what to OFFER; these decide what may happen, and they disagree on
|
||||
// purpose — a console rendered thirty seconds ago is a console describing a run
|
||||
// that has since moved.
|
||||
//
|
||||
// **A parked step is `running` with a NULL lease**, and that pair is the whole
|
||||
// vocabulary these need. `park()` above is the only thing that produces it, so
|
||||
// `status = 'running' AND claim_expires_at IS NULL` names a cue waiting on a
|
||||
// human and cannot name a step some process is mid-dispatch on. Confirm and skip
|
||||
// are both written against it, which is what makes them safe to expose to a
|
||||
// moderator: neither can touch a step the runner is holding.
|
||||
|
||||
/**
|
||||
* The highest `seq` of a step in this phase that is not still `pending` — the
|
||||
* furthest the phase has got — or null if none of it has been attempted.
|
||||
*
|
||||
* It exists for the retry control, and the definition is chosen to agree with
|
||||
* the runner's own cursor rather than to look tidy. Steps within a phase are
|
||||
* strictly serial, so the last step that is not pending is the last one the
|
||||
* runner worked on; if the run is `paused` that step is what it paused at.
|
||||
*
|
||||
* **The near miss worth recording: "the lowest step that is not settled" is the
|
||||
* wrong rule**, and it looks right. `nextOpenStep` selects `pending` and
|
||||
* `running` only, so a `failed` step is one the runner has already stepped OVER
|
||||
* — which is exactly what an `on_failure` of `skip` produces. Under that rule a
|
||||
* phase whose second step failed-and-skipped and whose fifth then failed-and-
|
||||
* paused would offer retry on the second, re-queueing a row behind the runner's
|
||||
* cursor where it would sit pending for ever.
|
||||
*/
|
||||
const lastStartedSeq = async (runId, phase) => {
|
||||
const [row] = await query(
|
||||
`SELECT MAX(seq) AS seq FROM event_run_steps
|
||||
WHERE run_id = ? AND phase = ? AND status <> 'pending'`,
|
||||
[runId, phase],
|
||||
)
|
||||
return row?.seq === null || row?.seq === undefined ? null : Number(row.seq)
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve a parked step: the GM cue's confirm.
|
||||
*
|
||||
* `done` rather than `skipped` — a human saying they did the thing is the step
|
||||
* having succeeded, and it is the only outcome under which the instruction was
|
||||
* actually carried out. The note is kept in `last_error` for the same reason the
|
||||
* park's is: it is the column the console already renders beside the step, and a
|
||||
* second one for prose would be a column two writers disagree about.
|
||||
*/
|
||||
const confirmParked = async (id, note) => {
|
||||
const result = await query(
|
||||
`UPDATE event_run_steps
|
||||
SET status = 'done', finished_at = NOW(), claimed_by = NULL,
|
||||
last_error = ?
|
||||
WHERE id = ? AND status = 'running' AND claim_expires_at IS NULL`,
|
||||
[note ? String(note).slice(0, 500) : null, id],
|
||||
)
|
||||
return Number(result?.affectedRows || 0) === 1
|
||||
}
|
||||
|
||||
/**
|
||||
* Skip a step a human has decided not to run: `pending`, or a parked cue.
|
||||
*
|
||||
* This is what `skipped` was reserved for (§L). A `running` step with a live
|
||||
* lease is excluded — nothing can recall a command already sent — and a `failed`
|
||||
* one is excluded because it is already terminal and the run's own resume is
|
||||
* what carries the phase past it.
|
||||
*/
|
||||
const skipByHuman = async (id, reason) => {
|
||||
const result = await query(
|
||||
`UPDATE event_run_steps
|
||||
SET status = 'skipped', finished_at = NOW(), claimed_by = NULL,
|
||||
last_error = ?
|
||||
WHERE id = ?
|
||||
AND (status = 'pending' OR (status = 'running' AND claim_expires_at IS NULL))`,
|
||||
[reason ? String(reason).slice(0, 500) : null, id],
|
||||
)
|
||||
return Number(result?.affectedRows || 0) === 1
|
||||
}
|
||||
|
||||
/**
|
||||
* Put a failed step back in the queue for another attempt.
|
||||
*
|
||||
* **`attempts` goes back to zero, and that is not the rule Engagement Phase 14
|
||||
* arrived at being broken.** That rule is about SWEEPS: an automatic path that
|
||||
* reset a counter made the ceiling unreachable and the row immortal. This is a
|
||||
* named person deciding, once, that the thing which failed three times will work
|
||||
* now — `EVENT_STEP_MAX_ATTEMPTS` bounds what the runner does unattended, and a
|
||||
* human is the thing it is unattended from. The decision is in the run log with
|
||||
* the actor on it.
|
||||
*/
|
||||
const requeue = async (id) => {
|
||||
const result = await query(
|
||||
`UPDATE event_run_steps
|
||||
SET status = 'pending', attempts = 0, due_at = NULL, last_error = NULL,
|
||||
claimed_by = NULL, claim_expires_at = NULL, finished_at = NULL
|
||||
WHERE id = ? AND status = 'failed'`,
|
||||
[id],
|
||||
)
|
||||
return Number(result?.affectedRows || 0) === 1
|
||||
}
|
||||
|
||||
/**
|
||||
* Close out every step a cancelled run will never run: pending, and parked.
|
||||
*
|
||||
* Wider than `cancelPending` by exactly one case, and deliberately so. §L leaves
|
||||
* a `running` step alone because nothing can recall a sent command — but a
|
||||
* parked cue is not a sent command, it is an instruction nobody is holding, and
|
||||
* leaving it `running` after the run was cancelled would leave the console
|
||||
* claiming a cancelled event is still waiting for someone. The live lease is
|
||||
* what distinguishes them, and it is in the WHERE clause.
|
||||
*/
|
||||
const cancelOpen = async (runId) => {
|
||||
const result = await query(
|
||||
`UPDATE event_run_steps
|
||||
SET status = 'cancelled', finished_at = NOW(), claimed_by = NULL
|
||||
WHERE run_id = ?
|
||||
AND (status = 'pending' OR (status = 'running' AND claim_expires_at IS NULL))`,
|
||||
[runId],
|
||||
)
|
||||
return Number(result?.affectedRows || 0)
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
listForRun,
|
||||
listForPhase,
|
||||
@@ -289,4 +413,9 @@ module.exports = {
|
||||
holdNext,
|
||||
reclaimStale,
|
||||
cancelPending,
|
||||
lastStartedSeq,
|
||||
confirmParked,
|
||||
skipByHuman,
|
||||
requeue,
|
||||
cancelOpen,
|
||||
}
|
||||
|
||||
@@ -23,8 +23,17 @@ const hydrate = (row) => row && { ...row, params: parseJson(row.params, null), r
|
||||
// only if every path a run can take reaches one of them.
|
||||
const TERMINAL = ['completed', 'cancelled', 'failed', 'missed']
|
||||
|
||||
// `waiting_steps` is the count of PARKED steps: `running` with a NULL lease, the
|
||||
// pair `park()` alone produces, which means a cue waiting on a human. It is a
|
||||
// correlated subquery on an admin list bounded at 500 rows rather than a column,
|
||||
// because it is derived from the steps and a column would be a second writer's
|
||||
// opinion of them. It earns its cost on the list screen: a cue nobody notices is
|
||||
// a run that never advances, and the run itself looks perfectly healthy until
|
||||
// somebody opens it.
|
||||
const SELECT_LIST = `
|
||||
SELECT r.*, d.title AS definition_title, d.slug AS definition_slug, v.version AS version_number
|
||||
SELECT r.*, d.title AS definition_title, d.slug AS definition_slug, v.version AS version_number,
|
||||
(SELECT COUNT(*) FROM event_run_steps s
|
||||
WHERE s.run_id = r.id AND s.status = 'running' AND s.claim_expires_at IS NULL) AS waiting_steps
|
||||
FROM event_runs r
|
||||
JOIN event_definitions d ON d.id = r.definition_id
|
||||
JOIN event_versions v ON v.id = r.version_id
|
||||
@@ -241,6 +250,20 @@ async function transition(id, from, to, { phase, error, clearClaim = false } = {
|
||||
return Number(result?.affectedRows || 0) === 1
|
||||
}
|
||||
|
||||
/**
|
||||
* Just this run's status, for a caller that must not act on a stale read.
|
||||
*
|
||||
* The runner drains a bounded batch of steps from one run inside a single tick,
|
||||
* and Phase 3 put a pause and a cancel button in a human's hand — so between two
|
||||
* steps of that batch the run may have stopped. A loop that only re-checked at
|
||||
* the top of the tick would answer a pause by dispatching another two dozen
|
||||
* steps, which is not a pause. One column, by primary key.
|
||||
*/
|
||||
const statusOf = async (id) => {
|
||||
const [row] = await query('SELECT status FROM event_runs WHERE id = ?', [id])
|
||||
return row?.status || null
|
||||
}
|
||||
|
||||
/**
|
||||
* Set health without touching status (§E).
|
||||
*
|
||||
@@ -354,6 +377,7 @@ module.exports = {
|
||||
claimStart,
|
||||
claimTick,
|
||||
releaseClaim,
|
||||
statusOf,
|
||||
transition,
|
||||
setHealth,
|
||||
concurrencyHolder,
|
||||
|
||||
@@ -8,10 +8,14 @@
|
||||
// a definition arriving from a future import or a restore gets the same answer
|
||||
// this screen does.
|
||||
//
|
||||
// **What is deliberately absent**: pause, resume, advance, cancel, step
|
||||
// skip/retry/confirm, cleanup and the action switchboard. Each of them acts on a
|
||||
// run in flight, and nothing is in flight until Phase 2 builds the runner. A
|
||||
// control that returns 200 and does nothing is worse than one that is not there.
|
||||
// **Phase 3 added the live run controls** at the bottom of this file: pause,
|
||||
// resume, cancel, and a step's confirm, skip and retry. What is still absent is
|
||||
// `advance`, `cleanup` and the action switchboard — `advance` has no honest
|
||||
// meaning until Phase 5 gives a phase an advance condition, `cleanup` has no
|
||||
// ledger to work over until Phase 8, and the switchboard is Phase 6's. Each of
|
||||
// them is absent rather than stubbed, for the reason the whole set was in Phase
|
||||
// 1: a control that returns 200 and does nothing is worse than one that is not
|
||||
// there.
|
||||
|
||||
const registries = require('../../../modules/registries')
|
||||
const spec = require('../../../events/spec')
|
||||
@@ -21,6 +25,7 @@ const versionsDb = require('../../../model/events/eventVersions.db')
|
||||
const seriesDb = require('../../../model/events/eventSeries.db')
|
||||
const runsDb = require('../../../model/events/eventRuns.db')
|
||||
const runs = require('../../../model/events/eventRuns.model')
|
||||
const controls = require('../../../model/events/eventRunControls.model')
|
||||
const logDb = require('../../../model/events/eventRunLog.db')
|
||||
const activity = require('../../../model/activity/activity.model')
|
||||
|
||||
@@ -80,6 +85,10 @@ const shapeRun = (r) => ({
|
||||
endedAt: r.ended_at,
|
||||
lastError: r.last_error,
|
||||
createdAt: r.created_at,
|
||||
// How many steps are parked on a human. Derived, not a column, and surfaced on
|
||||
// the LIST as well as the console because a cue nobody notices is a run that
|
||||
// never advances while looking perfectly healthy from the outside.
|
||||
waitingSteps: Number(r.waiting_steps || 0),
|
||||
})
|
||||
|
||||
const shapeStep = (s) => ({
|
||||
@@ -91,6 +100,12 @@ const shapeStep = (s) => ({
|
||||
params: s.params,
|
||||
actionVersion: s.action_version,
|
||||
status: s.status,
|
||||
// `running` with no lease is a parked step (§E) — waiting on a human, with
|
||||
// nothing holding it. The console has to tell that apart from a step some
|
||||
// process is mid-dispatch on, and it must not do so by being shown the lease:
|
||||
// one derived boolean rather than `claimed_by` and `claim_expires_at`, which
|
||||
// are the runner's business and would invite a UI that reasoned about leases.
|
||||
parked: s.status === 'running' && !s.claim_expires_at,
|
||||
dueAt: s.due_at,
|
||||
attempts: s.attempts,
|
||||
onFailure: s.on_failure,
|
||||
@@ -313,3 +328,89 @@ exports.startRun = async (req, res) => {
|
||||
created: result.created,
|
||||
})
|
||||
}
|
||||
|
||||
// ── Phase 3: the live run controls ─────────────────────────────────────────
|
||||
//
|
||||
// Six handlers, and each is the same four lines: read the ids out of the URL,
|
||||
// hand off to `eventRunControls`, log the manual transition to `activity_log`,
|
||||
// answer with the row. Every guard is in the model, where a control invoked from
|
||||
// anywhere else gets the same answer — which is the same division this file has
|
||||
// had since Phase 1.
|
||||
//
|
||||
// **The audit is written in two places on purpose, and they are not redundant.**
|
||||
// `event_run_log` is the run's own diagnostic record: queryable by phase and by
|
||||
// step, and it is what the console renders. `activity_log` is the deployment's
|
||||
// record of what staff did, and it is where "who cancelled the invasion" is
|
||||
// looked up months later by somebody who is not looking at that run. §J names
|
||||
// both.
|
||||
|
||||
/** POST /api/v1/admin/events/runs/:runId/pause */
|
||||
exports.pauseRun = async (req, res) => {
|
||||
const runId = asId(req.params.runId)
|
||||
if (!runId) return res.status(400).json({ error: 'bad run id' })
|
||||
const result = await controls.pause(runId, { reason: req.body?.reason }, req.user.id)
|
||||
if (!result.ok) return res.status(result.status || 409).json({ errors: result.errors })
|
||||
await activity.log({ req, action: 'event.run.paused', detail: { runId, reason: req.body?.reason || null } })
|
||||
return res.json({ run: shapeRun(result.run) })
|
||||
}
|
||||
|
||||
/** POST /api/v1/admin/events/runs/:runId/resume */
|
||||
exports.resumeRun = async (req, res) => {
|
||||
const runId = asId(req.params.runId)
|
||||
if (!runId) return res.status(400).json({ error: 'bad run id' })
|
||||
const result = await controls.resume(runId, {}, req.user.id)
|
||||
if (!result.ok) return res.status(result.status || 409).json({ errors: result.errors })
|
||||
await activity.log({ req, action: 'event.run.resumed', detail: { runId } })
|
||||
return res.json({ run: shapeRun(result.run) })
|
||||
}
|
||||
|
||||
/** POST /api/v1/admin/events/runs/:runId/cancel */
|
||||
exports.cancelRun = async (req, res) => {
|
||||
const runId = asId(req.params.runId)
|
||||
if (!runId) return res.status(400).json({ error: 'bad run id' })
|
||||
const result = await controls.cancel(runId, { reason: req.body?.reason }, req.user.id)
|
||||
if (!result.ok) return res.status(result.status || 409).json({ errors: result.errors })
|
||||
await activity.log({
|
||||
req,
|
||||
action: 'event.run.cancelled',
|
||||
detail: { runId, reason: req.body?.reason || null, cancelledSteps: result.cancelledSteps },
|
||||
})
|
||||
return res.json({ run: shapeRun(result.run), cancelledSteps: result.cancelledSteps })
|
||||
}
|
||||
|
||||
/** POST /api/v1/admin/events/runs/:runId/steps/:stepId/confirm */
|
||||
exports.confirmStep = async (req, res) => {
|
||||
const runId = asId(req.params.runId)
|
||||
const stepId = asId(req.params.stepId)
|
||||
if (!runId || !stepId) return res.status(400).json({ error: 'bad run or step id' })
|
||||
const result = await controls.confirmStep(runId, stepId, { note: req.body?.note }, req.user.id)
|
||||
if (!result.ok) return res.status(result.status || 409).json({ errors: result.errors })
|
||||
await activity.log({ req, action: 'event.step.confirmed', detail: { runId, stepId, action: result.step.action_id } })
|
||||
return res.json({ step: shapeStep(result.step) })
|
||||
}
|
||||
|
||||
/** POST /api/v1/admin/events/runs/:runId/steps/:stepId/skip */
|
||||
exports.skipStep = async (req, res) => {
|
||||
const runId = asId(req.params.runId)
|
||||
const stepId = asId(req.params.stepId)
|
||||
if (!runId || !stepId) return res.status(400).json({ error: 'bad run or step id' })
|
||||
const result = await controls.skipStep(runId, stepId, { reason: req.body?.reason }, req.user.id)
|
||||
if (!result.ok) return res.status(result.status || 409).json({ errors: result.errors })
|
||||
await activity.log({
|
||||
req,
|
||||
action: 'event.step.skipped',
|
||||
detail: { runId, stepId, action: result.step.action_id, reason: req.body?.reason || null },
|
||||
})
|
||||
return res.json({ step: shapeStep(result.step) })
|
||||
}
|
||||
|
||||
/** POST /api/v1/admin/events/runs/:runId/steps/:stepId/retry */
|
||||
exports.retryStep = async (req, res) => {
|
||||
const runId = asId(req.params.runId)
|
||||
const stepId = asId(req.params.stepId)
|
||||
if (!runId || !stepId) return res.status(400).json({ error: 'bad run or step id' })
|
||||
const result = await controls.retryStep(runId, stepId, {}, req.user.id)
|
||||
if (!result.ok) return res.status(result.status || 409).json({ errors: result.errors })
|
||||
await activity.log({ req, action: 'event.step.retried', detail: { runId, stepId, action: result.step.action_id } })
|
||||
return res.json({ step: shapeStep(result.step), run: shapeRun(result.run), resumed: result.resumed })
|
||||
}
|
||||
|
||||
@@ -12,9 +12,11 @@
|
||||
// are here anyway, because a button that is admin-only later and open now is a
|
||||
// gate nobody notices was missing.
|
||||
//
|
||||
// Reads are staff-wide. `verify` (admin, editor) is Phase 6's, and the live run
|
||||
// controls (admin, moderator) are Phase 3's — neither is stubbed here, because
|
||||
// nothing is in flight until Phase 2 builds the runner.
|
||||
// Reads are staff-wide. The live run controls landed in Phase 3 and are `admin`
|
||||
// + `moderator`, deliberately wider than start (§N2). `verify` (admin, editor),
|
||||
// `advance`, `cleanup` and the action switchboard are still absent rather than
|
||||
// stubbed — there is no advance condition until Phase 5, no resource ledger
|
||||
// until Phase 8 and no caps to price against until Phase 6.
|
||||
//
|
||||
// **Literal paths are declared before `/:id`**, so `/catalog`, `/series` and
|
||||
// `/runs` are never read as an event id.
|
||||
@@ -27,6 +29,10 @@ const { requireRole } = require('../../../utils/auth')
|
||||
const eventsRouter = express.Router()
|
||||
const adminOnly = requireRole('admin')
|
||||
const adminOrEditor = requireRole('admin', 'editor')
|
||||
// Live control of a run in flight, and the one gate wider than `admin` in this
|
||||
// feature (§K). Named rather than inlined so the six routes below cannot drift
|
||||
// apart from one another.
|
||||
const liveControl = requireRole('admin', 'moderator')
|
||||
|
||||
// ── The catalog and the vocabularies, served from the registries ───────────
|
||||
|
||||
@@ -93,6 +99,105 @@ eventsRouter.get(
|
||||
controller.getRunLog,
|
||||
)
|
||||
|
||||
// ── The live run controls (Phase 3) ───────────────────────────────────────
|
||||
//
|
||||
// `admin` + `moderator`, and it is the widest gate in this feature deliberately
|
||||
// (§K, §N2). Starting a run commits the deployment to everything the definition
|
||||
// contains, unattended, up to every cap it declares — that wants the narrowest
|
||||
// gate there is. Stopping one is incident response, and the incident is "the
|
||||
// event is doing something wrong at 2am" — that wants the widest. A split that
|
||||
// read consistent, with one role owning both buttons, would behave badly in
|
||||
// exactly the case the moderator role exists for.
|
||||
//
|
||||
// `advance` and `cleanup` from the § API surface table are not here: the first
|
||||
// has no honest meaning until Phase 5 gives a phase an advance condition, the
|
||||
// second has no resource ledger to work over until Phase 8.
|
||||
|
||||
eventsRouter.post(
|
||||
'/runs/:runId/pause',
|
||||
// #swagger.tags = ['Admin · Events']
|
||||
// #swagger.summary = 'Pause a run in flight'
|
||||
// #swagger.description = 'A paused run is excluded from the runner\'s sweep and nothing advances it until resume. Legal from `starting` and `running` only — a `scheduled` occurrence that should not happen is cancelled, not paused, because resuming one after its grace window had passed would produce a `missed` from a button labelled resume. Takes effect at once even mid-tick: the runner re-reads the run\'s status between steps.'
|
||||
// #swagger.security = [{ "cookieAuth": [] }, { "bearerAuth": [] }]
|
||||
/* #swagger.requestBody = { required: false, content: { "application/json": { schema: { type: "object", properties: { reason: { type: "string", description: "Recorded in the run log with the actor" } } } } } } */
|
||||
/* #swagger.responses[200] = { description: 'The paused run', content: { "application/json": { schema: { type: "object", properties: { run: { type: "object", additionalProperties: true } } } } } } */
|
||||
/* #swagger.responses[409] = { description: 'The run is not in flight', content: { "application/json": { schema: { type: "object", properties: { errors: { type: "array", items: { type: "string" } } } } } } } */
|
||||
/* #swagger.responses[403] = { description: 'Not an admin or moderator', content: { "application/json": { schema: { $ref: "#/components/schemas/Error" } } } } */
|
||||
liveControl,
|
||||
controller.pauseRun,
|
||||
)
|
||||
|
||||
eventsRouter.post(
|
||||
'/runs/:runId/resume',
|
||||
// #swagger.tags = ['Admin · Events']
|
||||
// #swagger.summary = 'Resume a paused run'
|
||||
// #swagger.description = 'Where the run goes back to is derived rather than remembered: a paused run with a `current_phase` was running, one without never got past `starting`. `last_error` is cleared — the operator has just dealt with it — and `health` is not, because "this run has already had trouble" stays true whoever pressed resume.'
|
||||
// #swagger.security = [{ "cookieAuth": [] }, { "bearerAuth": [] }]
|
||||
/* #swagger.responses[200] = { description: 'The resumed run', content: { "application/json": { schema: { type: "object", properties: { run: { type: "object", additionalProperties: true } } } } } } */
|
||||
/* #swagger.responses[409] = { description: 'The run is not paused', content: { "application/json": { schema: { type: "object", properties: { errors: { type: "array", items: { type: "string" } } } } } } } */
|
||||
/* #swagger.responses[403] = { description: 'Not an admin or moderator', content: { "application/json": { schema: { $ref: "#/components/schemas/Error" } } } } */
|
||||
liveControl,
|
||||
controller.resumeRun,
|
||||
)
|
||||
|
||||
eventsRouter.post(
|
||||
'/runs/:runId/cancel',
|
||||
// #swagger.tags = ['Admin · Events']
|
||||
// #swagger.summary = 'Cancel a run'
|
||||
// #swagger.description = 'Legal from every non-terminal status, `scheduled` included. Pending steps and any parked cue are cancelled with it; a step with a live lease is left alone, because nothing can recall a command already sent and a second writer on that row would race the process dispatching it. `cleanup` is not a parameter yet — the resource ledger it would work over arrives in Phase 8, and a flag that changes nothing is worse than one that is not there.'
|
||||
// #swagger.security = [{ "cookieAuth": [] }, { "bearerAuth": [] }]
|
||||
/* #swagger.requestBody = { required: false, content: { "application/json": { schema: { type: "object", properties: { reason: { type: "string", description: "Why. Recorded on the run and in its log, with the actor." } } } } } } */
|
||||
/* #swagger.responses[200] = { description: 'The cancelled run and how many steps were closed out with it', content: { "application/json": { schema: { type: "object", properties: { run: { type: "object", additionalProperties: true }, cancelledSteps: { type: "integer" } } } } } } */
|
||||
/* #swagger.responses[409] = { description: 'The run has already reached a terminal status', content: { "application/json": { schema: { type: "object", properties: { errors: { type: "array", items: { type: "string" } } } } } } } */
|
||||
/* #swagger.responses[403] = { description: 'Not an admin or moderator', content: { "application/json": { schema: { $ref: "#/components/schemas/Error" } } } } */
|
||||
liveControl,
|
||||
controller.cancelRun,
|
||||
)
|
||||
|
||||
eventsRouter.post(
|
||||
'/runs/:runId/steps/:stepId/confirm',
|
||||
// #swagger.tags = ['Admin · Events']
|
||||
// #swagger.summary = 'Confirm a parked step — the GM cue'
|
||||
// #swagger.description = 'The other half of `core.cue`. The action posts an instruction and parks the step `running` with a NULL lease — genuinely in flight, nothing holding it, so no sweep takes it back and a cue posted on Friday is still waiting on Monday. This ends it, as `done` rather than `skipped`: a person saying they did the thing is the step having succeeded. The optional note is what they did, and it is kept on the step and in the log.'
|
||||
// #swagger.security = [{ "cookieAuth": [] }, { "bearerAuth": [] }]
|
||||
/* #swagger.requestBody = { required: false, content: { "application/json": { schema: { type: "object", properties: { note: { type: "string", description: "What was actually done in-client" } } } } } } */
|
||||
/* #swagger.responses[200] = { description: 'The confirmed step', content: { "application/json": { schema: { type: "object", properties: { step: { type: "object", additionalProperties: true } } } } } } */
|
||||
/* #swagger.responses[404] = { description: 'No such run, or no such step on it', content: { "application/json": { schema: { $ref: "#/components/schemas/Error" } } } } */
|
||||
/* #swagger.responses[409] = { description: 'The step is not waiting on anyone', content: { "application/json": { schema: { type: "object", properties: { errors: { type: "array", items: { type: "string" } } } } } } } */
|
||||
/* #swagger.responses[403] = { description: 'Not an admin or moderator', content: { "application/json": { schema: { $ref: "#/components/schemas/Error" } } } } */
|
||||
liveControl,
|
||||
controller.confirmStep,
|
||||
)
|
||||
|
||||
eventsRouter.post(
|
||||
'/runs/:runId/steps/:stepId/skip',
|
||||
// #swagger.tags = ['Admin · Events']
|
||||
// #swagger.summary = 'Skip a step nobody is going to run'
|
||||
// #swagger.description = 'A step that has not started, or a parked cue. This is what the `skipped` status was reserved for, and why all three `on_failure` dispositions write `failed` instead — a status meaning both "a human decided against this" and "this was attempted three times and never worked" would make the console summary unreadable. A step with a live lease cannot be skipped; a failed one does not need to be, because resuming the run already carries the phase past it.'
|
||||
// #swagger.security = [{ "cookieAuth": [] }, { "bearerAuth": [] }]
|
||||
/* #swagger.requestBody = { required: false, content: { "application/json": { schema: { type: "object", properties: { reason: { type: "string" } } } } } } */
|
||||
/* #swagger.responses[200] = { description: 'The skipped step', content: { "application/json": { schema: { type: "object", properties: { step: { type: "object", additionalProperties: true } } } } } } */
|
||||
/* #swagger.responses[404] = { description: 'No such run, or no such step on it', content: { "application/json": { schema: { $ref: "#/components/schemas/Error" } } } } */
|
||||
/* #swagger.responses[409] = { description: 'The step or its run is in a status that cannot be skipped', content: { "application/json": { schema: { type: "object", properties: { errors: { type: "array", items: { type: "string" } } } } } } } */
|
||||
/* #swagger.responses[403] = { description: 'Not an admin or moderator', content: { "application/json": { schema: { $ref: "#/components/schemas/Error" } } } } */
|
||||
liveControl,
|
||||
controller.skipStep,
|
||||
)
|
||||
|
||||
eventsRouter.post(
|
||||
'/runs/:runId/steps/:stepId/retry',
|
||||
// #swagger.tags = ['Admin · Events']
|
||||
// #swagger.summary = 'Re-queue the failed step a paused run is stopped at, and resume it'
|
||||
// #swagger.description = 'One action rather than two, because there is no state in which you would want half of it: retry is legal only while the run is paused, and a paused run is paused AT this step. The step must be the one its phase is stopped at — a failed step under an `on_failure` of `skip` is one the run has already moved past, and re-queueing that would put a pending row behind the runner\'s cursor. `attempts` returns to zero: the attempt ceiling bounds what the runner does unattended, and a named person deciding is the thing it is unattended from.'
|
||||
// #swagger.security = [{ "cookieAuth": [] }, { "bearerAuth": [] }]
|
||||
/* #swagger.responses[200] = { description: 'The re-queued step and the run, with whether the resume took', content: { "application/json": { schema: { type: "object", properties: { step: { type: "object", additionalProperties: true }, run: { type: "object", additionalProperties: true }, resumed: { type: "boolean" } } } } } } */
|
||||
/* #swagger.responses[404] = { description: 'No such run, or no such step on it', content: { "application/json": { schema: { $ref: "#/components/schemas/Error" } } } } */
|
||||
/* #swagger.responses[409] = { description: 'The run is not paused, or the run is not stopped at this step', content: { "application/json": { schema: { type: "object", properties: { errors: { type: "array", items: { type: "string" } } } } } } } */
|
||||
/* #swagger.responses[403] = { description: 'Not an admin or moderator', content: { "application/json": { schema: { $ref: "#/components/schemas/Error" } } } } */
|
||||
liveControl,
|
||||
controller.retryStep,
|
||||
)
|
||||
|
||||
// ── Definitions ───────────────────────────────────────────────────────────
|
||||
|
||||
eventsRouter.get(
|
||||
|
||||
@@ -292,6 +292,14 @@ async function advanceRun(run, now) {
|
||||
const carry = {}
|
||||
|
||||
for (let n = 0; n < STEPS_PER_TICK; n++) {
|
||||
// Re-read the run's status between steps, not just at the top of the tick.
|
||||
// This loop drains up to STEPS_PER_TICK steps from one run, and Phase 3 put
|
||||
// a pause and a cancel in a human's hand: without this, pausing a run in the
|
||||
// middle of a batch would answer by dispatching another two dozen steps,
|
||||
// which is not a pause. One indexed column read per step, against a control
|
||||
// whose entire value is that it takes effect at once.
|
||||
if (n > 0 && (await runsDb.statusOf(run.id)) !== 'running') return 'stopped'
|
||||
|
||||
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` })
|
||||
|
||||
@@ -3766,6 +3766,101 @@
|
||||
]
|
||||
}
|
||||
},
|
||||
"/api/v1/admin/events/runs/{runId}/cancel": {
|
||||
"post": {
|
||||
"tags": [
|
||||
"Admin · Events"
|
||||
],
|
||||
"summary": "Cancel a run",
|
||||
"description": "Legal from every non-terminal status, `scheduled` included. Pending steps and any parked cue are cancelled with it; a step with a live lease is left alone, because nothing can recall a command already sent and a second writer on that row would race the process dispatching it. `cleanup` is not a parameter yet — the resource ledger it would work over arrives in Phase 8, and a flag that changes nothing is worse than one that is not there.",
|
||||
"parameters": [
|
||||
{
|
||||
"name": "runId",
|
||||
"in": "path",
|
||||
"required": true,
|
||||
"schema": {
|
||||
"type": "string"
|
||||
}
|
||||
}
|
||||
],
|
||||
"responses": {
|
||||
"200": {
|
||||
"description": "The cancelled run and how many steps were closed out with it",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"run": {
|
||||
"type": "object",
|
||||
"additionalProperties": true
|
||||
},
|
||||
"cancelledSteps": {
|
||||
"type": "integer"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"400": {
|
||||
"description": "Bad Request"
|
||||
},
|
||||
"403": {
|
||||
"description": "Not an admin or moderator",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/Error"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"409": {
|
||||
"description": "The run has already reached a terminal status",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"errors": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
"type": "string"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"security": [
|
||||
{
|
||||
"cookieAuth": []
|
||||
},
|
||||
{
|
||||
"bearerAuth": []
|
||||
}
|
||||
],
|
||||
"requestBody": {
|
||||
"required": false,
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"reason": {
|
||||
"type": "string",
|
||||
"description": "Why. Recorded on the run and in its log, with the actor."
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"/api/v1/admin/events/runs/{runId}/log": {
|
||||
"get": {
|
||||
"tags": [
|
||||
@@ -3842,6 +3937,494 @@
|
||||
]
|
||||
}
|
||||
},
|
||||
"/api/v1/admin/events/runs/{runId}/pause": {
|
||||
"post": {
|
||||
"tags": [
|
||||
"Admin · Events"
|
||||
],
|
||||
"summary": "Pause a run in flight",
|
||||
"description": "A paused run is excluded from the runner\\'s sweep and nothing advances it until resume. Legal from `starting` and `running` only — a `scheduled` occurrence that should not happen is cancelled, not paused, because resuming one after its grace window had passed would produce a `missed` from a button labelled resume. Takes effect at once even mid-tick: the runner re-reads the run\\'s status between steps.",
|
||||
"parameters": [
|
||||
{
|
||||
"name": "runId",
|
||||
"in": "path",
|
||||
"required": true,
|
||||
"schema": {
|
||||
"type": "string"
|
||||
}
|
||||
}
|
||||
],
|
||||
"responses": {
|
||||
"200": {
|
||||
"description": "The paused run",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"run": {
|
||||
"type": "object",
|
||||
"additionalProperties": true
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"400": {
|
||||
"description": "Bad Request"
|
||||
},
|
||||
"403": {
|
||||
"description": "Not an admin or moderator",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/Error"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"409": {
|
||||
"description": "The run is not in flight",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"errors": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
"type": "string"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"security": [
|
||||
{
|
||||
"cookieAuth": []
|
||||
},
|
||||
{
|
||||
"bearerAuth": []
|
||||
}
|
||||
],
|
||||
"requestBody": {
|
||||
"required": false,
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"reason": {
|
||||
"type": "string",
|
||||
"description": "Recorded in the run log with the actor"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"/api/v1/admin/events/runs/{runId}/resume": {
|
||||
"post": {
|
||||
"tags": [
|
||||
"Admin · Events"
|
||||
],
|
||||
"summary": "Resume a paused run",
|
||||
"description": "Where the run goes back to is derived rather than remembered: a paused run with a `current_phase` was running, one without never got past `starting`. `last_error` is cleared — the operator has just dealt with it — and `health` is not, because \"this run has already had trouble\" stays true whoever pressed resume.",
|
||||
"parameters": [
|
||||
{
|
||||
"name": "runId",
|
||||
"in": "path",
|
||||
"required": true,
|
||||
"schema": {
|
||||
"type": "string"
|
||||
}
|
||||
}
|
||||
],
|
||||
"responses": {
|
||||
"200": {
|
||||
"description": "The resumed run",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"run": {
|
||||
"type": "object",
|
||||
"additionalProperties": true
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"400": {
|
||||
"description": "Bad Request"
|
||||
},
|
||||
"403": {
|
||||
"description": "Not an admin or moderator",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/Error"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"409": {
|
||||
"description": "The run is not paused",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"errors": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
"type": "string"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"security": [
|
||||
{
|
||||
"cookieAuth": []
|
||||
},
|
||||
{
|
||||
"bearerAuth": []
|
||||
}
|
||||
]
|
||||
}
|
||||
},
|
||||
"/api/v1/admin/events/runs/{runId}/steps/{stepId}/confirm": {
|
||||
"post": {
|
||||
"tags": [
|
||||
"Admin · Events"
|
||||
],
|
||||
"summary": "Confirm a parked step — the GM cue",
|
||||
"description": "The other half of `core.cue`. The action posts an instruction and parks the step `running` with a NULL lease — genuinely in flight, nothing holding it, so no sweep takes it back and a cue posted on Friday is still waiting on Monday. This ends it, as `done` rather than `skipped`: a person saying they did the thing is the step having succeeded. The optional note is what they did, and it is kept on the step and in the log.",
|
||||
"parameters": [
|
||||
{
|
||||
"name": "runId",
|
||||
"in": "path",
|
||||
"required": true,
|
||||
"schema": {
|
||||
"type": "string"
|
||||
}
|
||||
},
|
||||
{
|
||||
"name": "stepId",
|
||||
"in": "path",
|
||||
"required": true,
|
||||
"schema": {
|
||||
"type": "string"
|
||||
}
|
||||
}
|
||||
],
|
||||
"responses": {
|
||||
"200": {
|
||||
"description": "The confirmed step",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"step": {
|
||||
"type": "object",
|
||||
"additionalProperties": true
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"400": {
|
||||
"description": "Bad Request"
|
||||
},
|
||||
"403": {
|
||||
"description": "Not an admin or moderator",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/Error"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"404": {
|
||||
"description": "No such run, or no such step on it",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/Error"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"409": {
|
||||
"description": "The step is not waiting on anyone",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"errors": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
"type": "string"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"security": [
|
||||
{
|
||||
"cookieAuth": []
|
||||
},
|
||||
{
|
||||
"bearerAuth": []
|
||||
}
|
||||
],
|
||||
"requestBody": {
|
||||
"required": false,
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"note": {
|
||||
"type": "string",
|
||||
"description": "What was actually done in-client"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"/api/v1/admin/events/runs/{runId}/steps/{stepId}/retry": {
|
||||
"post": {
|
||||
"tags": [
|
||||
"Admin · Events"
|
||||
],
|
||||
"summary": "Re-queue the failed step a paused run is stopped at, and resume it",
|
||||
"description": "One action rather than two, because there is no state in which you would want half of it: retry is legal only while the run is paused, and a paused run is paused AT this step. The step must be the one its phase is stopped at — a failed step under an `on_failure` of `skip` is one the run has already moved past, and re-queueing that would put a pending row behind the runner\\'s cursor. `attempts` returns to zero: the attempt ceiling bounds what the runner does unattended, and a named person deciding is the thing it is unattended from.",
|
||||
"parameters": [
|
||||
{
|
||||
"name": "runId",
|
||||
"in": "path",
|
||||
"required": true,
|
||||
"schema": {
|
||||
"type": "string"
|
||||
}
|
||||
},
|
||||
{
|
||||
"name": "stepId",
|
||||
"in": "path",
|
||||
"required": true,
|
||||
"schema": {
|
||||
"type": "string"
|
||||
}
|
||||
}
|
||||
],
|
||||
"responses": {
|
||||
"200": {
|
||||
"description": "The re-queued step and the run, with whether the resume took",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"step": {
|
||||
"type": "object",
|
||||
"additionalProperties": true
|
||||
},
|
||||
"run": {
|
||||
"type": "object",
|
||||
"additionalProperties": true
|
||||
},
|
||||
"resumed": {
|
||||
"type": "boolean"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"400": {
|
||||
"description": "Bad Request"
|
||||
},
|
||||
"403": {
|
||||
"description": "Not an admin or moderator",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/Error"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"404": {
|
||||
"description": "No such run, or no such step on it",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/Error"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"409": {
|
||||
"description": "The run is not paused, or the run is not stopped at this step",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"errors": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
"type": "string"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"security": [
|
||||
{
|
||||
"cookieAuth": []
|
||||
},
|
||||
{
|
||||
"bearerAuth": []
|
||||
}
|
||||
]
|
||||
}
|
||||
},
|
||||
"/api/v1/admin/events/runs/{runId}/steps/{stepId}/skip": {
|
||||
"post": {
|
||||
"tags": [
|
||||
"Admin · Events"
|
||||
],
|
||||
"summary": "Skip a step nobody is going to run",
|
||||
"description": "A step that has not started, or a parked cue. This is what the `skipped` status was reserved for, and why all three `on_failure` dispositions write `failed` instead — a status meaning both \"a human decided against this\" and \"this was attempted three times and never worked\" would make the console summary unreadable. A step with a live lease cannot be skipped; a failed one does not need to be, because resuming the run already carries the phase past it.",
|
||||
"parameters": [
|
||||
{
|
||||
"name": "runId",
|
||||
"in": "path",
|
||||
"required": true,
|
||||
"schema": {
|
||||
"type": "string"
|
||||
}
|
||||
},
|
||||
{
|
||||
"name": "stepId",
|
||||
"in": "path",
|
||||
"required": true,
|
||||
"schema": {
|
||||
"type": "string"
|
||||
}
|
||||
}
|
||||
],
|
||||
"responses": {
|
||||
"200": {
|
||||
"description": "The skipped step",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"step": {
|
||||
"type": "object",
|
||||
"additionalProperties": true
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"400": {
|
||||
"description": "Bad Request"
|
||||
},
|
||||
"403": {
|
||||
"description": "Not an admin or moderator",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/Error"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"404": {
|
||||
"description": "No such run, or no such step on it",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/Error"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"409": {
|
||||
"description": "The step or its run is in a status that cannot be skipped",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"errors": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
"type": "string"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"security": [
|
||||
{
|
||||
"cookieAuth": []
|
||||
},
|
||||
{
|
||||
"bearerAuth": []
|
||||
}
|
||||
],
|
||||
"requestBody": {
|
||||
"required": false,
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"reason": {
|
||||
"type": "string"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"/api/v1/admin/events/series": {
|
||||
"get": {
|
||||
"tags": [
|
||||
|
||||
442
server/test/eventRunControls.test.js
Normal file
442
server/test/eventRunControls.test.js
Normal file
@@ -0,0 +1,442 @@
|
||||
// ── The live run controls (EVENTS_PLAN.md Phase 3) ─────────────────────────
|
||||
//
|
||||
// Six controls, and what is tested is almost entirely the REFUSALS. A control
|
||||
// that works is easy; a control that works from a status it should not have
|
||||
// worked from is a staff member changing a live game world by pressing a button
|
||||
// a stale screen offered them. So each of the six is exercised from every status
|
||||
// it must decline, and the four that a run console could plausibly offer wrongly
|
||||
// get a test of their own:
|
||||
//
|
||||
// • retry on a step the run has already moved past (the `skip` disposition) —
|
||||
// the test that found the first draft's guard was reading the wrong end of
|
||||
// the phase
|
||||
// • confirm on a step a process is mid-dispatch on, not a parked cue
|
||||
// • skip on a step with a live lease
|
||||
// • cancel closing out a parked cue, so a cancelled run stops "waiting"
|
||||
//
|
||||
// The three tables are stubbed at the `.db` layer and the model's own logic runs
|
||||
// for real against them — the shape `eventRunner.test.js` uses. What a stub
|
||||
// cannot prove is that the five statements mean this against a real server; the
|
||||
// guards that are pure SQL (`status = 'running' AND claim_expires_at IS NULL`
|
||||
// and `lastStartedSeq`'s MAX) are proved in `eventRunnerSql.test.js`.
|
||||
//
|
||||
// Point the DB at a closed port before requiring anything.
|
||||
process.env.DB_HOST = '127.0.0.1'
|
||||
process.env.DB_PORT = '59999'
|
||||
|
||||
const { test, beforeEach, afterEach, after } = require('node:test')
|
||||
const assert = require('node:assert/strict')
|
||||
|
||||
const controls = require('../src/model/events/eventRunControls.model')
|
||||
const runsDb = require('../src/model/events/eventRuns.db')
|
||||
const stepsDb = require('../src/model/events/eventRunSteps.db')
|
||||
const logDb = require('../src/model/events/eventRunLog.db')
|
||||
const db = require('../src/utils/db')
|
||||
|
||||
after(() => db.close())
|
||||
|
||||
const TERMINAL = ['completed', 'cancelled', 'failed', 'missed']
|
||||
const ACTOR = 7
|
||||
|
||||
let store
|
||||
const originals = [
|
||||
['runs', runsDb, { ...runsDb }],
|
||||
['steps', stepsDb, { ...stepsDb }],
|
||||
['log', logDb, { ...logDb }],
|
||||
]
|
||||
|
||||
function installStubs() {
|
||||
store = { runs: new Map(), steps: new Map(), log: [], nextStepId: 1 }
|
||||
const snap = (o) => ({ ...o })
|
||||
|
||||
runsDb.getById = async (id) => {
|
||||
const r = store.runs.get(Number(id))
|
||||
return r ? snap(r) : null
|
||||
}
|
||||
|
||||
runsDb.transition = async (id, from, to, opts = {}) => {
|
||||
const r = store.runs.get(Number(id))
|
||||
const froms = Array.isArray(from) ? from : [from]
|
||||
if (!r || !froms.includes(r.status)) return false
|
||||
r.status = to
|
||||
if (opts.phase !== undefined) r.current_phase = opts.phase
|
||||
if (opts.error !== undefined) r.last_error = opts.error
|
||||
if (TERMINAL.includes(to) || opts.clearClaim) {
|
||||
r.claimed_by = null
|
||||
r.claim_expires_at = null
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
stepsDb.getById = async (id) => {
|
||||
const s = store.steps.get(Number(id))
|
||||
return s ? snap(s) : null
|
||||
}
|
||||
|
||||
// Each of the four mirrors its statement's WHERE clause exactly. A stub can
|
||||
// only ever agree with whoever wrote it, so what these buy is the model's
|
||||
// logic around them; the clauses themselves are checked against a real server.
|
||||
stepsDb.confirmParked = async (id, note) => {
|
||||
const s = store.steps.get(Number(id))
|
||||
if (!s || s.status !== 'running' || s.claim_expires_at) return false
|
||||
Object.assign(s, { status: 'done', last_error: note, claimed_by: null })
|
||||
return true
|
||||
}
|
||||
|
||||
stepsDb.skipByHuman = async (id, reason) => {
|
||||
const s = store.steps.get(Number(id))
|
||||
if (!s) return false
|
||||
const ok = s.status === 'pending' || (s.status === 'running' && !s.claim_expires_at)
|
||||
if (!ok) return false
|
||||
Object.assign(s, { status: 'skipped', last_error: reason, claimed_by: null })
|
||||
return true
|
||||
}
|
||||
|
||||
stepsDb.requeue = async (id) => {
|
||||
const s = store.steps.get(Number(id))
|
||||
if (!s || s.status !== 'failed') return false
|
||||
Object.assign(s, { status: 'pending', attempts: 0, due_at: null, last_error: null, claimed_by: null, claim_expires_at: null })
|
||||
return true
|
||||
}
|
||||
|
||||
stepsDb.cancelOpen = async (runId) => {
|
||||
let n = 0
|
||||
for (const s of store.steps.values()) {
|
||||
if (s.run_id !== Number(runId)) continue
|
||||
if (s.status === 'pending' || (s.status === 'running' && !s.claim_expires_at)) {
|
||||
s.status = 'cancelled'
|
||||
n += 1
|
||||
}
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
stepsDb.lastStartedSeq = async (runId, phase) => {
|
||||
const started = [...store.steps.values()]
|
||||
.filter((s) => s.run_id === Number(runId) && s.phase === phase && s.status !== 'pending')
|
||||
.map((s) => s.seq)
|
||||
return started.length ? Math.max(...started) : null
|
||||
}
|
||||
|
||||
logDb.write = async (line) => {
|
||||
store.log.push(line)
|
||||
return true
|
||||
}
|
||||
}
|
||||
|
||||
beforeEach(installStubs)
|
||||
afterEach(() => {
|
||||
for (const [, mod, fns] of originals) Object.assign(mod, fns)
|
||||
})
|
||||
|
||||
let nextRunId = 1
|
||||
|
||||
function seedRun({ status = 'running', phase = 'main', steps = [] } = {}) {
|
||||
const id = nextRunId++
|
||||
store.runs.set(id, {
|
||||
id,
|
||||
definition_id: id,
|
||||
version_id: id,
|
||||
status,
|
||||
health: 'ok',
|
||||
current_phase: phase,
|
||||
claimed_by: null,
|
||||
claim_expires_at: null,
|
||||
last_error: null,
|
||||
})
|
||||
steps.forEach((s, i) => {
|
||||
const stepId = store.nextStepId++
|
||||
store.steps.set(stepId, {
|
||||
id: stepId,
|
||||
run_id: id,
|
||||
phase: s.phase || phase,
|
||||
seq: s.seq ?? i,
|
||||
action_id: s.actionId || 'test.action',
|
||||
status: s.status || 'pending',
|
||||
attempts: s.attempts ?? 0,
|
||||
due_at: null,
|
||||
claimed_by: s.leased ? 'someone' : null,
|
||||
claim_expires_at: s.leased ? new Date(Date.now() + 60_000) : null,
|
||||
last_error: null,
|
||||
params: {},
|
||||
})
|
||||
})
|
||||
return id
|
||||
}
|
||||
|
||||
const runRow = (id) => store.runs.get(id)
|
||||
const stepsOf = (id) => [...store.steps.values()].filter((s) => s.run_id === id).sort((a, b) => a.seq - b.seq)
|
||||
const lastLog = () => store.log[store.log.length - 1]
|
||||
|
||||
// ── pause / resume ─────────────────────────────────────────────────────────
|
||||
|
||||
test('pause takes a run in flight and records who did it', async () => {
|
||||
const id = seedRun({ status: 'running' })
|
||||
const result = await controls.pause(id, { reason: 'the shard is lagging' }, ACTOR)
|
||||
|
||||
assert.equal(result.ok, true)
|
||||
assert.equal(runRow(id).status, 'paused')
|
||||
assert.deepEqual(lastLog().detail, {
|
||||
from: 'running',
|
||||
to: 'paused',
|
||||
control: 'pause',
|
||||
by: ACTOR,
|
||||
reason: 'the shard is lagging',
|
||||
})
|
||||
})
|
||||
|
||||
test('pause drops the claim, so the next tick is not locked out of a resumed run', async () => {
|
||||
const id = seedRun({ status: 'running' })
|
||||
Object.assign(runRow(id), { claimed_by: 'host:1', claim_expires_at: new Date(Date.now() + 900_000) })
|
||||
|
||||
await controls.pause(id, {}, ACTOR)
|
||||
|
||||
assert.equal(runRow(id).claimed_by, null)
|
||||
assert.equal(runRow(id).claim_expires_at, null)
|
||||
})
|
||||
|
||||
test('a scheduled run cannot be paused — it is cancelled instead', async () => {
|
||||
// Pausing one would leave a run that is neither going to start nor visibly
|
||||
// abandoned, and resuming it after its grace window had passed would produce a
|
||||
// `missed` from a button labelled resume.
|
||||
const id = seedRun({ status: 'scheduled' })
|
||||
const result = await controls.pause(id, {}, ACTOR)
|
||||
|
||||
assert.equal(result.ok, false)
|
||||
assert.equal(result.status, 409)
|
||||
assert.match(result.errors[0], /scheduled/)
|
||||
assert.equal(runRow(id).status, 'scheduled')
|
||||
})
|
||||
|
||||
test('a completed run cannot be paused', async () => {
|
||||
const id = seedRun({ status: 'completed' })
|
||||
assert.equal((await controls.pause(id, {}, ACTOR)).ok, false)
|
||||
})
|
||||
|
||||
test('resume returns a run to running, or to starting when it never entered a phase', async () => {
|
||||
const withPhase = seedRun({ status: 'paused', phase: 'main' })
|
||||
assert.equal((await controls.resume(withPhase, {}, ACTOR)).ok, true)
|
||||
assert.equal(runRow(withPhase).status, 'running')
|
||||
|
||||
const beforePhase = seedRun({ status: 'paused', phase: null })
|
||||
assert.equal((await controls.resume(beforePhase, {}, ACTOR)).ok, true)
|
||||
assert.equal(runRow(beforePhase).status, 'starting', 'both are in findDue; neither is a fourth column')
|
||||
})
|
||||
|
||||
test('resume clears the error it was paused over and leaves health alone', async () => {
|
||||
const id = seedRun({ status: 'paused' })
|
||||
Object.assign(runRow(id), { last_error: 'core.spawn failed', health: 'degraded' })
|
||||
|
||||
await controls.resume(id, {}, ACTOR)
|
||||
|
||||
assert.equal(runRow(id).last_error, null, 'a resolved failure must not accuse a healthy run for ever')
|
||||
assert.equal(runRow(id).health, 'degraded', 'that this run has already had trouble stays true')
|
||||
})
|
||||
|
||||
test('resume refuses a run that is not paused', async () => {
|
||||
const id = seedRun({ status: 'running' })
|
||||
const result = await controls.resume(id, {}, ACTOR)
|
||||
assert.equal(result.ok, false)
|
||||
assert.equal(result.status, 409)
|
||||
})
|
||||
|
||||
// ── cancel ─────────────────────────────────────────────────────────────────
|
||||
|
||||
test('cancel closes out the pending steps and the parked cue, and leaves a leased step alone', async () => {
|
||||
const id = seedRun({
|
||||
status: 'running',
|
||||
steps: [
|
||||
{ status: 'done' },
|
||||
{ status: 'running', leased: true }, // mid-dispatch: nothing can recall a sent command
|
||||
{ status: 'running' }, // parked on a human: nothing is holding it
|
||||
{ status: 'pending' },
|
||||
],
|
||||
})
|
||||
|
||||
const result = await controls.cancel(id, { reason: 'called off' }, ACTOR)
|
||||
|
||||
assert.equal(result.ok, true)
|
||||
assert.equal(runRow(id).status, 'cancelled')
|
||||
assert.equal(result.cancelledSteps, 2)
|
||||
const [done, leased, parked, pending] = stepsOf(id)
|
||||
assert.equal(done.status, 'done')
|
||||
assert.equal(leased.status, 'running', 'a step being dispatched is not touched')
|
||||
assert.equal(parked.status, 'cancelled', 'a cancelled run must stop claiming to wait on somebody')
|
||||
assert.equal(pending.status, 'cancelled')
|
||||
})
|
||||
|
||||
test('cancel is legal before a run has started', async () => {
|
||||
const id = seedRun({ status: 'scheduled', steps: [{ status: 'pending' }] })
|
||||
assert.equal((await controls.cancel(id, {}, ACTOR)).ok, true)
|
||||
assert.equal(runRow(id).status, 'cancelled')
|
||||
})
|
||||
|
||||
test('cancel refuses a run that is already terminal', async () => {
|
||||
for (const status of TERMINAL) {
|
||||
const id = seedRun({ status })
|
||||
const result = await controls.cancel(id, {}, ACTOR)
|
||||
assert.equal(result.ok, false, `${status} should not be cancellable`)
|
||||
assert.match(result.errors[0], new RegExp(status))
|
||||
}
|
||||
})
|
||||
|
||||
// ── confirm ────────────────────────────────────────────────────────────────
|
||||
|
||||
test('confirm resolves a parked cue as done, keeping what the person says they did', async () => {
|
||||
const id = seedRun({ status: 'running', steps: [{ status: 'running', actionId: 'core.cue' }] })
|
||||
const [cue] = stepsOf(id)
|
||||
|
||||
const result = await controls.confirmStep(id, cue.id, { note: 'gate opened, herald read' }, ACTOR)
|
||||
|
||||
assert.equal(result.ok, true)
|
||||
assert.equal(stepsOf(id)[0].status, 'done', 'a person saying they did it is the step having succeeded')
|
||||
assert.equal(stepsOf(id)[0].last_error, 'gate opened, herald read')
|
||||
assert.equal(lastLog().detail.control, 'confirm')
|
||||
assert.equal(lastLog().detail.by, ACTOR)
|
||||
})
|
||||
|
||||
test('confirm cannot resolve a step a process is dispatching', async () => {
|
||||
// The whole vocabulary here is "running with a NULL lease". A live lease means
|
||||
// something is mid-dispatch, and confirming it would race the process that
|
||||
// owns the row.
|
||||
const id = seedRun({ status: 'running', steps: [{ status: 'running', leased: true }] })
|
||||
const [busy] = stepsOf(id)
|
||||
|
||||
const result = await controls.confirmStep(id, busy.id, {}, ACTOR)
|
||||
|
||||
assert.equal(result.ok, false)
|
||||
assert.equal(result.status, 409)
|
||||
assert.equal(stepsOf(id)[0].status, 'running')
|
||||
})
|
||||
|
||||
test('a step id from another run is a 404, not an action', async () => {
|
||||
const mine = seedRun({ status: 'running', steps: [{ status: 'pending' }] })
|
||||
const theirs = seedRun({ status: 'running', steps: [{ status: 'running' }] })
|
||||
const [theirStep] = stepsOf(theirs)
|
||||
|
||||
const result = await controls.confirmStep(mine, theirStep.id, {}, ACTOR)
|
||||
|
||||
assert.equal(result.ok, false)
|
||||
assert.equal(result.status, 404)
|
||||
assert.equal(stepsOf(theirs)[0].status, 'running')
|
||||
})
|
||||
|
||||
// ── skip ───────────────────────────────────────────────────────────────────
|
||||
|
||||
test('skip takes a pending step and a parked cue, and nothing else', async () => {
|
||||
const id = seedRun({
|
||||
status: 'running',
|
||||
steps: [{ status: 'pending' }, { status: 'running' }, { status: 'running', leased: true }, { status: 'failed' }],
|
||||
})
|
||||
const [pending, parked, leased, failed] = stepsOf(id)
|
||||
|
||||
assert.equal((await controls.skipStep(id, pending.id, {}, ACTOR)).ok, true)
|
||||
assert.equal((await controls.skipStep(id, parked.id, {}, ACTOR)).ok, true)
|
||||
assert.equal((await controls.skipStep(id, leased.id, {}, ACTOR)).ok, false)
|
||||
// A failed step does not need skipping: `nextOpenStep` already passes over it,
|
||||
// so resuming the run carries the phase past it.
|
||||
assert.equal((await controls.skipStep(id, failed.id, {}, ACTOR)).ok, false)
|
||||
|
||||
const after = stepsOf(id)
|
||||
assert.equal(after[0].status, 'skipped')
|
||||
assert.equal(after[1].status, 'skipped')
|
||||
assert.equal(after[2].status, 'running')
|
||||
assert.equal(after[3].status, 'failed')
|
||||
})
|
||||
|
||||
test('skip refuses once the run is over', async () => {
|
||||
const id = seedRun({ status: 'completed', steps: [{ status: 'pending' }] })
|
||||
const [step] = stepsOf(id)
|
||||
assert.equal((await controls.skipStep(id, step.id, {}, ACTOR)).ok, false)
|
||||
})
|
||||
|
||||
// ── retry ──────────────────────────────────────────────────────────────────
|
||||
|
||||
test('retry re-queues the step a paused run is stopped at, and resumes in the same action', async () => {
|
||||
const id = seedRun({
|
||||
status: 'paused',
|
||||
steps: [{ status: 'done' }, { status: 'failed', attempts: 3 }, { status: 'pending' }],
|
||||
})
|
||||
const failed = stepsOf(id)[1]
|
||||
|
||||
const result = await controls.retryStep(id, failed.id, {}, ACTOR)
|
||||
|
||||
assert.equal(result.ok, true)
|
||||
assert.equal(result.resumed, true)
|
||||
assert.equal(stepsOf(id)[1].status, 'pending')
|
||||
assert.equal(stepsOf(id)[1].attempts, 0, 'the ceiling bounds the runner, not a person deciding once')
|
||||
assert.equal(runRow(id).status, 'running', 'there is no state in which you would want half of this')
|
||||
})
|
||||
|
||||
test('retry refuses a step the run has already moved past', async () => {
|
||||
// The case the guard exists for: a failed step under an `on_failure` of `skip`
|
||||
// is one the phase carried on from. Re-queueing it would put a pending row
|
||||
// behind the runner's cursor, where it would sit for ever.
|
||||
const id = seedRun({
|
||||
status: 'paused',
|
||||
steps: [{ status: 'failed', attempts: 3 }, { status: 'done' }, { status: 'failed', attempts: 3 }],
|
||||
})
|
||||
const [movedPast] = stepsOf(id)
|
||||
|
||||
const result = await controls.retryStep(id, movedPast.id, {}, ACTOR)
|
||||
|
||||
assert.equal(result.ok, false)
|
||||
assert.equal(result.status, 409)
|
||||
assert.match(result.errors[0], /stopped at this step/)
|
||||
assert.equal(stepsOf(id)[0].status, 'failed')
|
||||
assert.equal(runRow(id).status, 'paused', 'a refused retry does not resume the run either')
|
||||
})
|
||||
|
||||
test('retry refuses a step in a phase the run has left', async () => {
|
||||
const id = seedRun({
|
||||
status: 'paused',
|
||||
phase: 'two',
|
||||
steps: [{ phase: 'one', seq: 0, status: 'failed' }, { phase: 'two', seq: 0, status: 'pending' }],
|
||||
})
|
||||
const [old] = stepsOf(id)
|
||||
|
||||
const result = await controls.retryStep(id, old.id, {}, ACTOR)
|
||||
assert.equal(result.ok, false)
|
||||
assert.match(result.errors[0], /already left/)
|
||||
})
|
||||
|
||||
test('retry refuses while the run is still running', async () => {
|
||||
const id = seedRun({ status: 'running', steps: [{ status: 'failed' }] })
|
||||
const [failed] = stepsOf(id)
|
||||
|
||||
const result = await controls.retryStep(id, failed.id, {}, ACTOR)
|
||||
assert.equal(result.ok, false)
|
||||
assert.match(result.errors[0], /paused/)
|
||||
})
|
||||
|
||||
test('retry refuses a step that is not failed', async () => {
|
||||
const id = seedRun({ status: 'paused', steps: [{ status: 'pending' }] })
|
||||
const [pending] = stepsOf(id)
|
||||
assert.equal((await controls.retryStep(id, pending.id, {}, ACTOR)).ok, false)
|
||||
})
|
||||
|
||||
// ── the record ─────────────────────────────────────────────────────────────
|
||||
|
||||
test('every control writes one log line carrying the actor and the control name', async () => {
|
||||
const id = seedRun({ status: 'running', steps: [{ status: 'running' }, { status: 'pending' }] })
|
||||
const [parked, pending] = stepsOf(id)
|
||||
|
||||
await controls.confirmStep(id, parked.id, { note: 'done' }, ACTOR)
|
||||
await controls.skipStep(id, pending.id, { reason: 'not needed' }, ACTOR)
|
||||
await controls.pause(id, {}, ACTOR)
|
||||
await controls.resume(id, {}, ACTOR)
|
||||
await controls.cancel(id, { reason: 'over' }, ACTOR)
|
||||
|
||||
const human = store.log.filter((l) => l.detail?.control)
|
||||
assert.deepEqual(human.map((l) => l.detail.control), ['confirm', 'skip', 'pause', 'resume', 'cancel'])
|
||||
assert.ok(human.every((l) => l.detail.by === ACTOR))
|
||||
// The kinds are the ones a reader already scans for. A human transition is
|
||||
// still a transition; `detail.control` is what separates it from the runner's.
|
||||
assert.deepEqual([...new Set(human.map((l) => l.kind))].sort(), ['run.status', 'step.status'])
|
||||
})
|
||||
|
||||
test('an empty reason is stored as NULL rather than as an empty string', async () => {
|
||||
const id = seedRun({ status: 'running' })
|
||||
await controls.pause(id, { reason: ' ' }, ACTOR)
|
||||
assert.equal(lastLog().detail.reason, null)
|
||||
})
|
||||
@@ -131,6 +131,8 @@ function installStubs() {
|
||||
return true
|
||||
}
|
||||
|
||||
runsDb.statusOf = async (id) => store.runs.get(id)?.status || null
|
||||
|
||||
runsDb.setHealth = async (id, health) => {
|
||||
const r = store.runs.get(id)
|
||||
if (!r || r.health === health) return false
|
||||
@@ -687,3 +689,51 @@ test('the dispatch envelope carries what §F says it carries', async () => {
|
||||
assert.equal(envelope.idempotencyKey.length, 40)
|
||||
assert.deepEqual(Object.keys(envelope).sort(), ['actor', 'idempotencyKey', 'params', 'runId', 'scope', 'stepId', 'verify'])
|
||||
})
|
||||
|
||||
test('a run paused mid-batch stops there rather than draining the rest of the phase', async () => {
|
||||
// The whole value of a pause is that it takes effect NOW. `advanceRun` drains
|
||||
// up to STEPS_PER_TICK steps from one run inside a single tick, so a status
|
||||
// re-read only at the top of the tick would answer a pause by dispatching
|
||||
// another two dozen steps. Written by pausing from inside an action's own
|
||||
// `perform`, which is the only moment that race is reproducible.
|
||||
register([
|
||||
scriptedAction('test.pauser', {
|
||||
perform: async ({ runId }) => {
|
||||
await runsDb.transition(runId, ['starting', 'running'], 'paused', { clearClaim: true })
|
||||
return { ok: true }
|
||||
},
|
||||
}),
|
||||
scriptedAction('test.after'),
|
||||
])
|
||||
|
||||
const id = seedRun([
|
||||
{
|
||||
key: 'main',
|
||||
label: 'Main',
|
||||
steps: [step('test.pauser'), step('test.after'), step('test.after')],
|
||||
},
|
||||
])
|
||||
|
||||
await runner.tick(T0)
|
||||
|
||||
assert.equal(run(id).status, 'paused')
|
||||
const [first, second, third] = stepsOf(id)
|
||||
assert.equal(first.status, 'done', 'the step that was already dispatched finishes')
|
||||
assert.equal(second.status, 'pending', 'nothing after it ran')
|
||||
assert.equal(third.status, 'pending')
|
||||
assert.equal(scripted['test.after'], undefined, 'the later action was never called')
|
||||
})
|
||||
|
||||
test('a resumed run picks up from the step it stopped at', async () => {
|
||||
register([scriptedAction('test.a'), scriptedAction('test.b')])
|
||||
|
||||
const id = seedRun([{ key: 'main', label: 'Main', steps: [step('test.a'), step('test.b')] }])
|
||||
await runner.tick(T0)
|
||||
assert.equal(run(id).status, 'completed')
|
||||
|
||||
// And the mirror of it: a run parked at `paused` is not picked up at all, which
|
||||
// is what `findDue`'s omission of the status buys.
|
||||
const other = seedRun([{ key: 'main', label: 'Main', steps: [step('test.a')] }], { status: 'paused' })
|
||||
await runner.tick(T0)
|
||||
assert.equal(stepsOf(other)[0].status, 'pending', 'a paused run is not swept')
|
||||
})
|
||||
|
||||
@@ -25,6 +25,19 @@
|
||||
// both a MariaDB-specific syntax and a correctness claim: it must move the
|
||||
// next PENDING step and only ever push a due date later.
|
||||
//
|
||||
// **Phase 3 added four more**, and each of them is a control a staff member
|
||||
// presses against a live game world:
|
||||
//
|
||||
// * **`confirmParked` / `skipByHuman`** - both keyed on `status = 'running'
|
||||
// AND claim_expires_at IS NULL`. That pair, and only that pair, means "a cue
|
||||
// waiting on a human". If the clause let a LEASED step through, confirm would
|
||||
// race the process mid-dispatch on that row.
|
||||
// * **`cancelOpen`** - pending steps and parked cues, never a leased one.
|
||||
// * **`lastStartedSeq`** - a `MAX(seq) ... WHERE status <> 'pending'`, which is
|
||||
// what decides whether retry is offered. The first draft asked for the LOWEST
|
||||
// unsettled seq instead, which is a different step whenever a phase carried
|
||||
// on past an `on_failure: skip` failure.
|
||||
//
|
||||
// Plus the two unique indexes that are load-bearing rather than tidy:
|
||||
// `uq_evrun_occurrence` (which, not the claim, is what stops two runs of one
|
||||
// occurrence existing) and `uq_evstep_slot` (which is what makes re-materialising
|
||||
@@ -150,6 +163,36 @@ UPDATE event_run_steps
|
||||
AND (due_at IS NULL OR due_at < ?)
|
||||
ORDER BY seq LIMIT 1`
|
||||
|
||||
// Phase 3's four, verbatim from `eventRunSteps.db.js`.
|
||||
const CONFIRM_PARKED = `
|
||||
UPDATE event_run_steps
|
||||
SET status = 'done', finished_at = NOW(), claimed_by = NULL,
|
||||
last_error = ?
|
||||
WHERE id = ? AND status = 'running' AND claim_expires_at IS NULL`
|
||||
|
||||
const SKIP_BY_HUMAN = `
|
||||
UPDATE event_run_steps
|
||||
SET status = 'skipped', finished_at = NOW(), claimed_by = NULL,
|
||||
last_error = ?
|
||||
WHERE id = ?
|
||||
AND (status = 'pending' OR (status = 'running' AND claim_expires_at IS NULL))`
|
||||
|
||||
const REQUEUE = `
|
||||
UPDATE event_run_steps
|
||||
SET status = 'pending', attempts = 0, due_at = NULL, last_error = NULL,
|
||||
claimed_by = NULL, claim_expires_at = NULL, finished_at = NULL
|
||||
WHERE id = ? AND status = 'failed'`
|
||||
|
||||
const CANCEL_OPEN = `
|
||||
UPDATE event_run_steps
|
||||
SET status = 'cancelled', finished_at = NOW(), claimed_by = NULL
|
||||
WHERE run_id = ?
|
||||
AND (status = 'pending' OR (status = 'running' AND claim_expires_at IS NULL))`
|
||||
|
||||
const LAST_STARTED_SEQ = `
|
||||
SELECT MAX(seq) AS seq FROM event_run_steps
|
||||
WHERE run_id = ? AND phase = ? AND status <> 'pending'`
|
||||
|
||||
const MATERIALISE_RUN = `
|
||||
INSERT IGNORE INTO event_runs (definition_id, version_id, scope, scheduled_for, concurrency_key)
|
||||
VALUES (?, ?, ?, ?, ?)`
|
||||
@@ -506,3 +549,113 @@ test('findMissed compares against each definition’s own grace window', async (
|
||||
assert.deepEqual(missed, [tight.runId])
|
||||
assert.ok(!missed.includes(generous.runId), 'inside its own window a run starts late rather than being missed')
|
||||
})
|
||||
|
||||
// -- Phase 3: the controls a human presses ----------------------------------
|
||||
|
||||
test('confirm resolves a parked cue and cannot touch a step being dispatched', async (t) => {
|
||||
if (needDb(t)) return
|
||||
const { runId } = await seedRun({ status: 'running' })
|
||||
const parked = await seedStep(runId, { seq: 0, status: 'running', claimedBy: 'host:1', claimExpiresAt: null })
|
||||
const busy = await seedStep(runId, { seq: 1, status: 'running', claimedBy: 'host:1', claimExpiresAt: later(60_000), key: 'x'.repeat(40) })
|
||||
|
||||
assert.equal(rows(await pool.query(CONFIRM_PARKED, ['gate opened', parked])), 1)
|
||||
assert.equal(rows(await pool.query(CONFIRM_PARKED, ['nope', busy])), 0, 'a live lease is a step somebody owns')
|
||||
|
||||
assert.equal((await stepById(parked)).status, 'done')
|
||||
assert.equal((await stepById(parked)).last_error, 'gate opened')
|
||||
assert.equal((await stepById(busy)).status, 'running')
|
||||
})
|
||||
|
||||
test('a confirm of an already-confirmed cue reports 0, not 1', async (t) => {
|
||||
if (needDb(t)) return
|
||||
// The engagement Phase 4a shape: a connector that defaults `foundRows: true`
|
||||
// reports 1 for an UPDATE that matched and changed nothing, and a control that
|
||||
// read that as success would tell a second staff member their press worked.
|
||||
const { runId } = await seedRun({ status: 'running' })
|
||||
const parked = await seedStep(runId, { status: 'running', claimExpiresAt: null })
|
||||
|
||||
assert.equal(rows(await pool.query(CONFIRM_PARKED, [null, parked])), 1)
|
||||
assert.equal(rows(await pool.query(CONFIRM_PARKED, [null, parked])), 0)
|
||||
})
|
||||
|
||||
test('skip takes a pending step and a parked cue, and refuses a leased one', async (t) => {
|
||||
if (needDb(t)) return
|
||||
const { runId } = await seedRun({ status: 'running' })
|
||||
const pending = await seedStep(runId, { seq: 0, status: 'pending' })
|
||||
const parked = await seedStep(runId, { seq: 1, status: 'running', claimExpiresAt: null, key: 'y'.repeat(40) })
|
||||
const busy = await seedStep(runId, { seq: 2, status: 'running', claimExpiresAt: later(60_000), key: 'z'.repeat(40) })
|
||||
const failed = await seedStep(runId, { seq: 3, status: 'failed', key: 'w'.repeat(40) })
|
||||
|
||||
assert.equal(rows(await pool.query(SKIP_BY_HUMAN, [null, pending])), 1)
|
||||
assert.equal(rows(await pool.query(SKIP_BY_HUMAN, [null, parked])), 1)
|
||||
assert.equal(rows(await pool.query(SKIP_BY_HUMAN, [null, busy])), 0)
|
||||
assert.equal(rows(await pool.query(SKIP_BY_HUMAN, [null, failed])), 0, 'a failed step is terminal; resume carries the phase past it')
|
||||
})
|
||||
|
||||
test('cancelOpen closes pending steps and parked cues, and leaves a leased one alone', async (t) => {
|
||||
if (needDb(t)) return
|
||||
const { runId } = await seedRun({ status: 'running' })
|
||||
const done = await seedStep(runId, { seq: 0, status: 'done' })
|
||||
const busy = await seedStep(runId, { seq: 1, status: 'running', claimExpiresAt: later(60_000), key: 'p'.repeat(40) })
|
||||
const parked = await seedStep(runId, { seq: 2, status: 'running', claimExpiresAt: null, key: 'q'.repeat(40) })
|
||||
const pending = await seedStep(runId, { seq: 3, status: 'pending', key: 'r'.repeat(40) })
|
||||
|
||||
assert.equal(rows(await pool.query(CANCEL_OPEN, [runId])), 2)
|
||||
|
||||
assert.equal((await stepById(done)).status, 'done')
|
||||
assert.equal((await stepById(busy)).status, 'running', 'nothing can recall a command already sent')
|
||||
assert.equal((await stepById(parked)).status, 'cancelled', 'a cancelled run must stop claiming to wait on somebody')
|
||||
assert.equal((await stepById(pending)).status, 'cancelled')
|
||||
})
|
||||
|
||||
test('requeue only takes a failed step, and puts attempts back to zero', async (t) => {
|
||||
if (needDb(t)) return
|
||||
const { runId } = await seedRun({ status: 'paused' })
|
||||
const failed = await seedStep(runId, { seq: 0, status: 'failed', attempts: 3 })
|
||||
const pending = await seedStep(runId, { seq: 1, status: 'pending', key: 's'.repeat(40) })
|
||||
|
||||
assert.equal(rows(await pool.query(REQUEUE, [failed])), 1)
|
||||
assert.equal(rows(await pool.query(REQUEUE, [pending])), 0)
|
||||
|
||||
const row = await stepById(failed)
|
||||
assert.equal(row.status, 'pending')
|
||||
assert.equal(Number(row.attempts), 0)
|
||||
assert.equal(row.due_at, null, 'a re-queued step is due now, not at the retry backoff it was left on')
|
||||
})
|
||||
|
||||
test('lastStartedSeq names the furthest step of the phase, not the earliest unsettled one', async (t) => {
|
||||
if (needDb(t)) return
|
||||
// The defect this replaced: a phase that carried on past a failed step (an
|
||||
// `on_failure` of `skip`) and then paused at a later one. "The lowest seq that
|
||||
// is not settled" answers with the FIRST failure - a step the runner has long
|
||||
// since stepped over - and retry would re-queue a row behind its own cursor.
|
||||
const { runId } = await seedRun({ status: 'paused' })
|
||||
await seedStep(runId, { seq: 0, status: 'failed', key: 'a'.repeat(40) })
|
||||
await seedStep(runId, { seq: 1, status: 'done', key: 'b'.repeat(40) })
|
||||
await seedStep(runId, { seq: 2, status: 'failed', key: 'c'.repeat(40) })
|
||||
await seedStep(runId, { seq: 3, status: 'pending', key: 'd'.repeat(40) })
|
||||
|
||||
const [row] = await pool.query(LAST_STARTED_SEQ, [runId, 'main'])
|
||||
assert.equal(Number(row.seq), 2)
|
||||
})
|
||||
|
||||
test('lastStartedSeq is NULL for a phase nothing has touched', async (t) => {
|
||||
if (needDb(t)) return
|
||||
const { runId } = await seedRun({ status: 'running' })
|
||||
await seedStep(runId, { seq: 0, status: 'pending' })
|
||||
|
||||
const [row] = await pool.query(LAST_STARTED_SEQ, [runId, 'main'])
|
||||
assert.equal(row.seq, null, 'a null must read as "nothing to retry", not as seq 0')
|
||||
})
|
||||
|
||||
test('a guarded transition refuses a run that was cancelled underneath it', async (t) => {
|
||||
if (needDb(t)) return
|
||||
// What the admin cancel looks like from the runner's side, mid-tick: the
|
||||
// guarded write returns 0 and the tick treats the run as taken rather than
|
||||
// advancing a run somebody has just stopped.
|
||||
const { runId } = await seedRun({ status: 'running' })
|
||||
await pool.query("UPDATE event_runs SET status = 'cancelled' WHERE id = ?", [runId])
|
||||
|
||||
assert.equal(rows(await pool.query(TRANSITION, ['running', 'two', runId, 'running'])), 0)
|
||||
assert.equal((await runById(runId)).status, 'cancelled')
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user