const { query } = require('../../utils/db') const { parseJson } = require('./engagementRules.db') const hydrate = (row) => row && { ...row, payload: parseJson(row.payload, {}) } /** * Enqueue one (rule, user, channel) row, idempotently. * * `INSERT IGNORE` rather than a plain INSERT, because `uq_engo_dedupe` is the * replay guard (§4.2a): the sidecar feed is at-least-once and a reconnect * backfills, so the same event arriving twice must produce one row and not two * mails. IGNORE turns that into a silent no-op, which is what a replay should be. * * Returns the new id, or null when the row already existed. A null is a * SUCCESSFUL duplicate, not a failure - the caller counts it as such. * * A NULL dedupe_key never collides (multiple NULLs are legal under a UNIQUE * index), so an emit that carries no key always enqueues. That is the right * default: dedupe is something the emitter opts into by naming a key, and core * cannot invent one that means anything. */ async function enqueue(row) { const result = await query( `INSERT IGNORE INTO engagement_outbox (rule_id, trigger_id, user_id, channel, subject_key, scope_key, payload, dedupe_key, due_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`, [ row.rule_id, row.trigger_id, row.user_id, row.channel, row.subject_key || '', // NULL, not '', for an unscoped event: '' is a scope key that means // "deployment-wide" in engagement_digest_state, and this column has to be // able to say "no scope at all" as well. row.scope_key ?? null, JSON.stringify(row.payload || {}), row.dedupe_key ?? null, row.due_at, ], ) return Number(result?.affectedRows || 0) === 1 ? result.insertId : null } /** * Rows that are due. `idx_engo_due (status, due_at)` is this query. * * It selects rather than claims - claiming is `claim()` below, one row at a * time - so two instances sweeping at once both see the same candidates and then * disagree, harmlessly, about which of them owns each. */ const findDue = async (now, limit = 100) => ( await query( "SELECT * FROM engagement_outbox WHERE status = 'scheduled' AND due_at <= ? ORDER BY due_at, id LIMIT ?", [now, limit], ) ).map(hydrate) /** * Take ownership of one due row: a compare-and-set from 'scheduled' to 'sending'. * * **This is §7.1 Q2's answer** (settled by the org lead 2026-08-29, over * `SELECT ... FOR UPDATE SKIP LOCKED`). The winner is whoever the server reports * `affectedRows = 1` to; every other sweeper gets 0 and moves on. No explicit * transaction, no MariaDB version floor, and it uses a status the ENUM already * carried for exactly this. * * What it makes safe is the OUTBOX and only the outbox. `announceWorker`, * `teamDigestWorker`, `teamForumUploadSweep` and `teamActivityPrune` are all * still written for a single instance, so this does not make the deployment * multi-instance - it makes the one table that will carry mail ready for the day * it is, which is cheap now and expensive after mail has doubled once. */ async function claim(id) { const result = await query( `UPDATE engagement_outbox SET status = 'sending', attempts = attempts + 1 WHERE id = ? AND status = 'scheduled'`, [id], ) return Number(result?.affectedRows || 0) === 1 } /** * Release a claimed row back to 'scheduled' with a later `due_at` - a transient * failure that should be retried. The mirror of announceJobs' backoff. */ const reschedule = (id, dueAt, error) => query( "UPDATE engagement_outbox SET status = 'scheduled', due_at = ?, last_error = ? WHERE id = ? AND status = 'sending'", [dueAt, error ? String(error).slice(0, 2000) : null, id], ) /** A terminal outcome: 'sent', 'failed' or 'suppressed'. */ const finish = (id, status, error) => query( `UPDATE engagement_outbox SET status = ?, last_error = ?, sent_at = IF(? = 'sent', NOW(), sent_at) WHERE id = ?`, [status, error ? String(error).slice(0, 2000) : null, status, id], ) /** * Cancel every still-scheduled row for a (rule, subject) - the point of the * grace window (§4.2a). `userId` narrows it to one recipient when the resolving * event names one; a resolving event with no owner cancels for everyone the * original event was queued for, which is the house-repaired case. * * Only 'scheduled' rows are touched: a row already claimed into 'sending' is * somebody's in-flight send and cancelling it would leave two workers writing * one row's outcome. */ async function cancel(ruleId, subjectKey, userId = null) { const params = [ruleId, subjectKey] let sql = "UPDATE engagement_outbox SET status = 'cancelled' WHERE rule_id = ? AND subject_key = ? AND status = 'scheduled'" if (userId !== null && userId !== undefined) { sql += ' AND user_id = ?' params.push(userId) } const result = await query(sql, params) return Number(result?.affectedRows || 0) } /** * Recover rows stranded in 'sending' by a crash between the claim and the * outcome. * * Without this the CAS claim leaks: the claiming process died, no other sweeper * will ever match `status = 'scheduled'`, and the row sits in 'sending' forever. * `updated_at` is the clock (it is ON UPDATE CURRENT_TIMESTAMP, so the claim * stamped it), and the window has to be comfortably longer than the slowest * legitimate send or this reclaims rows that are merely slow. */ const reclaimStale = async (before, maxAttempts = 0) => { // Give up first, reclaim second, and in that order: a row that has already // burned its attempts must leave 'sending' as `failed`, or the reclaim below // hands it straight back to `findDue` and it is retried forever. // // **This is what makes `pruneTerminal` a bound at all** (Phase 14). Attempts // are incremented by `claim`, but `MAX_ATTEMPTS` is only consulted on a // graceful `retry` outcome — a send that kills the process mid-flight never // reaches that branch, so before this the row cycled sending → scheduled → // sending forever, never reached a terminal status, and was therefore never // eligible for any retention sweep. One poisoned payload was an outbox row // that outlived every horizon. let failed = 0 if (Number(maxAttempts) > 0) { const gaveUp = await query( `UPDATE engagement_outbox SET status = 'failed', last_error = 'gave up after repeated interruptions' WHERE status = 'sending' AND updated_at < ? AND attempts >= ?`, [before, Math.floor(maxAttempts)], ) failed = Number(gaveUp?.affectedRows || 0) } const reclaimed = await query( "UPDATE engagement_outbox SET status = 'scheduled' WHERE status = 'sending' AND updated_at < ?", [before], ) return { failed, reclaimed: Number(reclaimed?.affectedRows || 0) } } /** * Delete terminal rows older than `before` (Phase 14). * * **Terminal only, and the status list is the whole policy.** A `scheduled` row * is a promise the engine has not kept yet — `delay_seconds` can legitimately * put one up to a day out (`MAX_DELAY_SECONDS`) — and a `sending` row may be a * worker mid-flight. Deleting either is not retention, it is cancelling a send * nobody asked to cancel. Only `sent`, `failed`, `cancelled` and `suppressed` * are outcomes that have already happened. * * `created_at` rather than `updated_at` is the clock deliberately: the horizon * an operator sets means "how long we keep the record of a delivery", which is * measured from when it was enqueued, not from whenever it was last touched. * * @returns {Promise} rows deleted */ const pruneTerminal = async (before, limit = 1000) => { const result = await query( `DELETE FROM engagement_outbox WHERE status IN ('sent', 'failed', 'cancelled', 'suppressed') AND created_at < ? LIMIT ?`, [before, Math.floor(limit)], ) return Number(result?.affectedRows || 0) } const getById = async (id) => { const [row] = await query('SELECT * FROM engagement_outbox WHERE id = ?', [id]) return hydrate(row) } /** Admin/read surfaces (Phase 4b) and tests. */ const listForRule = async (ruleId, limit = 100) => ( await query('SELECT * FROM engagement_outbox WHERE rule_id = ? ORDER BY id DESC LIMIT ?', [ruleId, limit]) ).map(hydrate) module.exports = { enqueue, findDue, claim, reschedule, finish, cancel, reclaimStale, pruneTerminal, getById, listForRule, }