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
422 lines
17 KiB
JavaScript
422 lines
17 KiB
JavaScript
// ── event_run_steps — SQL only ─────────────────────────────────────────────
|
|
//
|
|
// EVENTS.md §D and §E. Phase 1 materialises a run's steps and reads them back
|
|
// for the run console. **Draining them is Phase 2's**: the CAS claim, the lease,
|
|
// the attempt counter and the classification of a module's answer are the
|
|
// runner, and none of them is stubbed here.
|
|
//
|
|
// The one runtime property Phase 1 does have to get right is the idempotency key
|
|
// (§E). Core mints it ONCE, at materialisation, and it does NOT vary by attempt —
|
|
// a retry re-sends the same key so the game side can recognise the repeat. That
|
|
// makes it a property of the INSERT below rather than of the dispatch, which is
|
|
// the only reason it can be stable at all.
|
|
|
|
const crypto = require('crypto')
|
|
|
|
const { query } = require('../../utils/db')
|
|
const { parseJson } = require('./eventJson')
|
|
|
|
const hydrate = (row) => row && { ...row, params: parseJson(row.params, {}) }
|
|
|
|
/**
|
|
* `sha256(runId | stepId)`, truncated to 40 hex — the shape `shardEvents.dedupeKey`
|
|
* already uses, so the two dedupe keys on this codebase read alike.
|
|
*
|
|
* The step id is not known until the row exists, so materialisation inserts with
|
|
* a provisional key and stamps the real one immediately afterwards. That is one
|
|
* extra statement per step and it buys the property the whole retry story rests
|
|
* on: the key is a function of identity, never of attempt or of clock.
|
|
*/
|
|
const idempotencyKey = (runId, stepId) =>
|
|
crypto.createHash('sha256').update(`${runId}|${stepId}`).digest('hex').slice(0, 40)
|
|
|
|
const listForRun = async (runId) =>
|
|
(
|
|
await query(
|
|
'SELECT * FROM event_run_steps WHERE run_id = ? ORDER BY phase, seq, id',
|
|
[runId],
|
|
)
|
|
).map(hydrate)
|
|
|
|
const getById = async (id) => {
|
|
const [row] = await query('SELECT * FROM event_run_steps WHERE id = ?', [id])
|
|
return hydrate(row)
|
|
}
|
|
|
|
/**
|
|
* Materialise one phase's steps.
|
|
*
|
|
* `INSERT IGNORE` against `UNIQUE (run_id, phase, seq)`, so a tick that overran
|
|
* into the next one cannot double-materialise a phase — the same argument the
|
|
* occurrence key makes one table up, at the other end of the run.
|
|
*
|
|
* Returns the rows as they now stand, created or pre-existing, so a caller that
|
|
* lost the race still gets the step ids.
|
|
*/
|
|
const materialisePhase = async (runId, phase, steps) => {
|
|
for (let i = 0; i < steps.length; i++) {
|
|
const step = steps[i]
|
|
const result = await query(
|
|
`INSERT IGNORE INTO event_run_steps
|
|
(run_id, phase, seq, action_id, params, action_version, on_failure, idempotency_key)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, '')`,
|
|
[
|
|
runId,
|
|
phase,
|
|
i,
|
|
step.actionId,
|
|
JSON.stringify(step.params || {}),
|
|
step.actionVersion || 1,
|
|
step.onFailure || 'pause',
|
|
],
|
|
)
|
|
if (Number(result?.affectedRows || 0) === 1) {
|
|
// Stamped in a second statement because the key is a function of the row's
|
|
// own id. Scoped by the empty key so a re-run of this loop over an existing
|
|
// phase can never overwrite a key a dispatch has already sent.
|
|
await query(
|
|
"UPDATE event_run_steps SET idempotency_key = ? WHERE id = ? AND idempotency_key = ''",
|
|
[idempotencyKey(runId, result.insertId), result.insertId],
|
|
)
|
|
}
|
|
}
|
|
return listForRun(runId)
|
|
}
|
|
|
|
/** The run console's summary line: how many steps sit in each status. */
|
|
const statusCounts = async (runId) => {
|
|
const rows = await query(
|
|
'SELECT status, COUNT(*) AS n FROM event_run_steps WHERE run_id = ? GROUP BY status',
|
|
[runId],
|
|
)
|
|
return Object.fromEntries(rows.map((r) => [r.status, Number(r.n)]))
|
|
}
|
|
|
|
// ── Phase 2: draining a step ───────────────────────────────────────────────
|
|
//
|
|
// The CAS claim, the lease, the attempt counter and the terminal writes. Phase 1
|
|
// left all of it out rather than stubbing it, and this is where it lands.
|
|
//
|
|
// **Two rules govern everything below, and both were paid for once already.**
|
|
//
|
|
// 1. `attempts` is incremented by the CLAIM and by nothing else, and no recovery
|
|
// path ever resets it. Engagement Phase 14's defect was a stale-row sweep that
|
|
// returned rows to their start state: the attempt ceiling became unreachable,
|
|
// so the row cycled forever, never reached a terminal status, and was
|
|
// therefore never eligible for any retention sweep.
|
|
// 2. A PARKED step is `running` with a NULL lease, and the reclaim only ever
|
|
// touches a lease that is non-NULL and expired (the org lead's answer,
|
|
// 2026-09-02). That is what lets a GM cue wait for a human overnight without a
|
|
// sweep re-dispatching the instruction every fifteen minutes.
|
|
|
|
// A run's steps in authored order, for the phase the run is currently in.
|
|
const listForPhase = async (runId, phase) =>
|
|
(
|
|
await query(
|
|
'SELECT * FROM event_run_steps WHERE run_id = ? AND phase = ? ORDER BY seq, id',
|
|
[runId, phase],
|
|
)
|
|
).map(hydrate)
|
|
|
|
/**
|
|
* The next step of a phase that the runner may work on, or null.
|
|
*
|
|
* **Steps within a phase are strictly serial.** This returns the lowest-`seq`
|
|
* step that is not terminal, and the runner does nothing with step N+1 until N
|
|
* has finished — which is the only reading under which `core.wait` means anything
|
|
* at all, and the only one under which a cue can gate what follows it.
|
|
*
|
|
* A parked or running step is returned too, so the caller can see that the phase
|
|
* is occupied rather than concluding it is finished.
|
|
*/
|
|
const nextOpenStep = async (runId, phase) => {
|
|
const [row] = await query(
|
|
`SELECT * FROM event_run_steps
|
|
WHERE run_id = ? AND phase = ?
|
|
AND status IN ('pending','running')
|
|
ORDER BY seq, id LIMIT 1`,
|
|
[runId, phase],
|
|
)
|
|
return hydrate(row) || null
|
|
}
|
|
|
|
/**
|
|
* Take ownership of one pending step: the CAS `pending -> running`, plus a lease.
|
|
*
|
|
* `due_at` is honoured here rather than in the caller's filter so that the whole
|
|
* decision — is it mine, is it due — is one statement the database arbitrates. A
|
|
* NULL `due_at` is due now, which is what materialisation writes for every step
|
|
* that is not sitting behind a `core.wait`.
|
|
*/
|
|
async function claim(id, owner, leaseUntil, now) {
|
|
const result = await query(
|
|
`UPDATE event_run_steps
|
|
SET status = 'running', attempts = attempts + 1, claimed_by = ?, claim_expires_at = ?,
|
|
started_at = COALESCE(started_at, NOW())
|
|
WHERE id = ? AND status = 'pending' AND (due_at IS NULL OR due_at <= ?)`,
|
|
[owner, leaseUntil, id, now],
|
|
)
|
|
return Number(result?.affectedRows || 0) === 1
|
|
}
|
|
|
|
/**
|
|
* Park a claimed step: it stays `running`, and its lease goes NULL.
|
|
*
|
|
* This is the whole mechanism behind `core.cue`. The step is genuinely in flight
|
|
* — an instruction has been posted and nothing else in the phase may proceed —
|
|
* but no process is holding it, so the reclaim must not take it back. A NULL
|
|
* lease says exactly that, and `reclaimStale` below is written to agree.
|
|
*/
|
|
const park = (id, note) =>
|
|
query(
|
|
`UPDATE event_run_steps SET claim_expires_at = NULL, last_error = ?
|
|
WHERE id = ? AND status = 'running'`,
|
|
[note ? String(note).slice(0, 500) : null, id],
|
|
)
|
|
|
|
/**
|
|
* Release a claimed step back to `pending` for a later attempt.
|
|
*
|
|
* `attempts` is untouched — it was already incremented by the claim, which is the
|
|
* only place that may. Backoff is flat rather than exponential for the reason the
|
|
* outbox's is: `due_at` is also the event's own clock, and a doubling backoff
|
|
* pushes a step arbitrarily far past the moment the event was about.
|
|
*/
|
|
const reschedule = (id, dueAt, error) =>
|
|
query(
|
|
`UPDATE event_run_steps
|
|
SET status = 'pending', due_at = ?, claimed_by = NULL, claim_expires_at = NULL, last_error = ?
|
|
WHERE id = ? AND status = 'running'`,
|
|
[dueAt, error ? String(error).slice(0, 500) : null, id],
|
|
)
|
|
|
|
/** A terminal outcome for one step: done, failed, skipped, refused or cancelled. */
|
|
const finish = (id, status, error) =>
|
|
query(
|
|
`UPDATE event_run_steps
|
|
SET status = ?, last_error = ?, finished_at = NOW(),
|
|
claimed_by = NULL, claim_expires_at = NULL
|
|
WHERE id = ? AND status = 'running'`,
|
|
[status, error ? String(error).slice(0, 500) : null, id],
|
|
)
|
|
|
|
/**
|
|
* Delay the next not-yet-started step of a phase — what `core.wait` actually does.
|
|
*
|
|
* The wait step itself completes normally; the pause is the NEXT step's `due_at`,
|
|
* owned by the runner. A `perform()` that slept would hold its claim for the
|
|
* duration and turn a five-minute pause into a five-minute lease, which is the
|
|
* one shape this must not have.
|
|
*
|
|
* Guarded on `status = 'pending'` and on the current `due_at` being sooner, so a
|
|
* re-dispatch of a wait whose ack was lost cannot push the following step further
|
|
* out a second time.
|
|
*/
|
|
const holdNext = async (runId, phase, afterSeq, dueAt) => {
|
|
const result = await query(
|
|
`UPDATE event_run_steps
|
|
SET due_at = ?
|
|
WHERE run_id = ? AND phase = ? AND seq > ? AND status = 'pending'
|
|
AND (due_at IS NULL OR due_at < ?)
|
|
ORDER BY seq LIMIT 1`,
|
|
[dueAt, runId, phase, afterSeq, dueAt],
|
|
)
|
|
return Number(result?.affectedRows || 0) === 1
|
|
}
|
|
|
|
/**
|
|
* Recover steps whose claim outlived the process that took it.
|
|
*
|
|
* **`attempts` is not reset and the lease being NULL is not staleness.** The
|
|
* first is Engagement Phase 14's rule; the second is what makes a parked cue
|
|
* survive. A step that has already burned its attempts leaves `running` as
|
|
* `failed` rather than being handed back, and in that order — a reclaim that ran
|
|
* first would return it to `pending` and it would be retried forever.
|
|
*/
|
|
const reclaimStale = async (now, maxAttempts = 0) => {
|
|
let failed = 0
|
|
if (Number(maxAttempts) > 0) {
|
|
const gaveUp = await query(
|
|
`UPDATE event_run_steps
|
|
SET status = 'failed', last_error = 'gave up after repeated interruptions',
|
|
finished_at = NOW(), claimed_by = NULL, claim_expires_at = NULL
|
|
WHERE status = 'running'
|
|
AND claim_expires_at IS NOT NULL AND claim_expires_at < ?
|
|
AND attempts >= ?`,
|
|
[now, Math.floor(maxAttempts)],
|
|
)
|
|
failed = Number(gaveUp?.affectedRows || 0)
|
|
}
|
|
const reclaimed = await query(
|
|
`UPDATE event_run_steps
|
|
SET status = 'pending', claimed_by = NULL, claim_expires_at = NULL
|
|
WHERE status = 'running' AND claim_expires_at IS NOT NULL AND claim_expires_at < ?`,
|
|
[now],
|
|
)
|
|
return { failed, reclaimed: Number(reclaimed?.affectedRows || 0) }
|
|
}
|
|
|
|
/**
|
|
* Cancel every step of a run that has not started (Phase 3's cancel, and the
|
|
* abort_run disposition).
|
|
*
|
|
* A `running` step is deliberately left alone, parked or not: nothing can recall
|
|
* a command already sent, and a second writer on that row would race the process
|
|
* that owns it (§L).
|
|
*/
|
|
const cancelPending = async (runId) => {
|
|
const result = await query(
|
|
`UPDATE event_run_steps
|
|
SET status = 'cancelled', finished_at = NOW()
|
|
WHERE run_id = ? AND status = 'pending'`,
|
|
[runId],
|
|
)
|
|
return Number(result?.affectedRows || 0)
|
|
}
|
|
|
|
// ── 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,
|
|
getById,
|
|
materialisePhase,
|
|
statusCounts,
|
|
idempotencyKey,
|
|
nextOpenStep,
|
|
claim,
|
|
park,
|
|
reschedule,
|
|
finish,
|
|
holdNext,
|
|
reclaimStale,
|
|
cancelPending,
|
|
lastStartedSeq,
|
|
confirmParked,
|
|
skipByHuman,
|
|
requeue,
|
|
cancelOpen,
|
|
}
|