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, payload, dedupe_key, due_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?)`, [ row.rule_id, row.trigger_id, row.user_id, row.channel, row.subject_key || '', 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 = (before) => query( "UPDATE engagement_outbox SET status = 'scheduled' WHERE status = 'sending' AND updated_at < ?", [before], ) 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, getById, listForRule, }