Files
website/server/test/engagementEngine.test.js
wtclaude 065bec7ad8
All checks were successful
PR Checks / bot-tests (pull_request) Successful in 28s
PR Checks / client-build (pull_request) Successful in 29s
PR Checks / server-tests (pull_request) Successful in 11m9s
feat(engagement): the email channel on the engine, and the Teams migration (engagement Phase 6)
Email becomes a DeliveryChannel driven by rules, and the Team pipeline stops being
its own thing. `teamNotify.forumPost` now emits an event; a rule decides who is
mailed, through which template, and how often at most. One walk goes forum write
-> events.emit -> rule -> outbox -> worker -> email channel -> template -> SMTP.

Seven decisions settled by the org lead before any code:

  - email only moves; the push tickle and the Discord bridge stay direct calls
  - the EVENT carries its access-checked audience, and `members` resolves to it
  - the four Team rules are seeded DISABLED, with an admin banner and a note
  - team_notification_prefs stays, read by the engine as a scoped preference
  - the payload wins and a structural projection fills the gaps
  - the digest keeps computing at send time; only its state generalizes
  - an unsubscribe token turns off the channel it names, and nothing else

Three defects found while building it:

  - `email.button` never absolutized its href, while image and itemList both
    did. Every rule-driven CTA would have been a dead relative link, because a
    trigger's url variables are validated site-relative by construction.
  - Phase 4a enqueued digest-mode recipients for a drain that Phase 6 decided
    not to build. An outbox row snapshots the payload and so has none of the
    three properties the digest design exists for, including the security one.
  - the digest's send-log row carried no address_hash while the instant row
    beside it did, which would have made half the mail uncorrelatable in Phase 9.

Also: engagement_digest_state + a replay-safe backfill, engagement_outbox.scope_key,
a v2 unsubscribe token that still verifies v1 forever, and the canonical
/public/engagement/unsubscribe pair with the old /public/teams path kept
permanently — mail is not editable once sent.

Verified with 1464 server tests, 324 client tests, and a live rig (MariaDB +
Mailpit + a real Team) covering the instant mail, the digest, the generic
template, a pre-migration unsubscribe link and the backfill's replay-safety.

Docs: RunicGateway/docs#TBD

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-29 20:11:54 -05:00

916 lines
36 KiB
JavaScript

// ── The engagement engine (ENGAGEMENT.md Phase 4a) ─────────────────────────
//
// The phase's acceptance criteria, one test apiece:
//
// • a trigger fired twice inside `cooldown_seconds` for the same
// (rule, user, subject) sends once
// • the same trigger for a DIFFERENT subject sends again — the multi-house
// case §4.1 names, which is the one a per-user cooldown gets wrong
// • a scheduled row is cancelled by a `cancel_on` trigger and never sends
// • a restart mid-window still sends exactly once
// • a duplicate `dedupe_key` is a successful no-op
//
// …plus the two properties that are security boundaries rather than behaviour:
// the G24 ceiling is re-checked at SEND time and not only at save, and a composed
// segment takes the NARROWEST ceiling in its tree.
//
// **The five tables are stubbed at the `.db` layer** and the engine's own logic
// runs for real against them, the shape `notificationChannelPrefs.test.js` uses.
// The one place that is not enough is the raw SQL whose correctness IS a server
// contract - the cooldown claim, the outbox compare-and-set, and the scoped
// dedupe key. Those run against a real MariaDB in `engagementEngineSql.test.js`,
// which skips when there is none, and the first time it ran it disproved the
// cooldown statement this file's stub had been agreeing with.
//
// Point the DB at a closed port before requiring anything: the registries reach
// utils/discordAnnounce, which builds the pool at require time.
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 registries = require('../src/modules/registries')
const channels = require('../src/engagement/channels')
const engine = require('../src/engagement/engine')
const conditions = require('../src/engagement/conditions')
const segments = require('../src/engagement/segments')
const worker = require('../src/utils/engagementWorker')
const rules = require('../src/model/engagement/engagementRules.model')
const rulesDb = require('../src/model/engagement/engagementRules.db')
const outboxDb = require('../src/model/engagement/engagementOutbox.db')
const cooldownsDb = require('../src/model/engagement/engagementCooldowns.db')
const sendsDb = require('../src/model/engagement/engagementSends.db')
const segmentsDb = require('../src/model/engagement/engagementSegments.db')
const recipients = require('../src/model/engagement/engagementRecipients.db')
const db = require('../src/utils/db')
after(() => db.close())
const T0 = new Date('2026-08-29T12:00:00Z')
const later = (ms) => new Date(T0.getTime() + ms)
// ── In-memory stand-ins for the five tables ────────────────────────────────
let store
const originals = {}
function snapshotOriginals() {
for (const [name, mod] of [
['rulesDb', rulesDb], ['outboxDb', outboxDb], ['cooldownsDb', cooldownsDb],
['sendsDb', sendsDb], ['segmentsDb', segmentsDb], ['recipients', recipients],
]) {
originals[name] = { mod, fns: { ...mod } }
}
}
snapshotOriginals()
function restoreOriginals() {
for (const { mod, fns } of Object.values(originals)) Object.assign(mod, fns)
}
function installStubs() {
store = {
rules: new Map(),
segments: new Map(),
cooldowns: new Map(),
outbox: new Map(),
sends: [],
users: new Map(), // id -> { id, role, status }
prefs: new Map(), // "<user> <id> <channel>" -> mode
nextOutboxId: 1,
}
rulesDb.enabledForTrigger = async (triggerId) =>
[...store.rules.values()].filter((r) => r.enabled && r.trigger_id === triggerId)
rulesDb.enabledCancelledBy = async (triggerId) =>
[...store.rules.values()].filter((r) => r.enabled && (r.cancel_on || []).includes(triggerId))
rulesDb.getById = async (id) => store.rules.get(id) || null
rulesDb.list = async () => [...store.rules.values()]
rulesDb.countUsingSegment = async (segmentId) =>
[...store.rules.values()].filter((r) => r.audience_segment_id === segmentId).length
segmentsDb.getById = async (id) => store.segments.get(id) || null
segmentsDb.list = async () => [...store.segments.values()]
// The two statements' semantics, reproduced: a guarded UPDATE that matches
// claims the fire; otherwise an INSERT IGNORE claims a first fire; otherwise
// the pair is still cooling. `engagementEngineSql.test.js` is what proves the
// SQL itself - a stub can only ever agree with whoever wrote it, and in this
// case the first version of both was wrong together.
cooldownsDb.claim = async (ruleId, userId, subjectKey, cooldownSeconds, now) => {
const key = `${ruleId}|${userId}|${subjectKey}`
const row = store.cooldowns.get(key)
if (!row) {
store.cooldowns.set(key, { last_fired_at: now, fire_count: 1 })
return true
}
if (row.last_fired_at.getTime() <= now.getTime() - cooldownSeconds * 1000) {
row.fire_count += 1
row.last_fired_at = now
return true
}
return false
}
outboxDb.enqueue = async (row) => {
if (row.dedupe_key) {
const clash = [...store.outbox.values()].find(
(r) =>
r.dedupe_key === row.dedupe_key &&
r.rule_id === row.rule_id &&
r.user_id === row.user_id &&
r.channel === row.channel,
)
if (clash) return null
}
const id = store.nextOutboxId++
store.outbox.set(id, { id, status: 'scheduled', attempts: 0, subject_key: '', ...row })
return id
}
// Copies, not the live objects: a SQL SELECT hands back a snapshot, and
// `processRow` reads `row.attempts` as the value BEFORE its own claim
// incremented it. Returning references here made the retry budget off by one
// in the stub only, which is exactly the class of thing a stub must not invent.
outboxDb.findDue = async (now, limit = 100) =>
[...store.outbox.values()]
.filter((r) => r.status === 'scheduled' && r.due_at <= now)
.sort((a, b) => a.due_at - b.due_at || a.id - b.id)
.slice(0, limit)
.map((r) => ({ ...r }))
outboxDb.claim = async (id) => {
const row = store.outbox.get(id)
if (!row || row.status !== 'scheduled') return false
row.status = 'sending'
row.attempts += 1
return true
}
outboxDb.reschedule = async (id, dueAt, error) => {
const row = store.outbox.get(id)
if (row && row.status === 'sending') Object.assign(row, { status: 'scheduled', due_at: dueAt, last_error: error })
}
outboxDb.finish = async (id, status, error) => {
const row = store.outbox.get(id)
if (row) Object.assign(row, { status, last_error: error })
}
outboxDb.cancel = async (ruleId, subjectKey, userId = null) => {
let n = 0
for (const row of store.outbox.values()) {
if (row.rule_id !== ruleId || row.subject_key !== subjectKey || row.status !== 'scheduled') continue
if (userId !== null && userId !== undefined && row.user_id !== userId) continue
row.status = 'cancelled'
n += 1
}
return n
}
outboxDb.reclaimStale = async () => {}
outboxDb.getById = async (id) => (store.outbox.has(id) ? { ...store.outbox.get(id) } : null)
sendsDb.record = async (entry) => {
store.sends.push({ id: store.sends.length + 1, created_at: T0, ...entry })
return store.sends.length
}
sendsDb.countSentSince = async (ruleId, since) =>
store.sends.filter((s) => s.rule_id === ruleId && s.status === 'sent' && s.created_at >= since).length
const activeIds = () => [...store.users.values()].filter((u) => u.status === 'active').map((u) => u.id)
recipients.active = async () => activeIds()
recipients.staff = async (roles) =>
[...store.users.values()].filter((u) => u.status === 'active' && roles.includes(u.role)).map((u) => u.id)
recipients.subscribers = async (streamId, defaultOn = []) =>
activeIds().filter((id) => {
const rows = [...store.prefs.entries()].filter(([k]) => k.startsWith(`${id} ${streamId} `))
if (rows.some(([, mode]) => mode !== 'off')) return true
const named = new Set(rows.map(([k]) => k.split(' ')[2]))
return defaultOn.some((c) => !named.has(c))
})
recipients.filterActive = async (ids) =>
[...new Set(ids)].filter((id) => store.users.get(id)?.status === 'active')
recipients.storedModes = async (userIds, streamId, channel) =>
new Map(
userIds
.filter((id) => store.prefs.has(`${id} ${streamId} ${channel}`))
.map((id) => [id, store.prefs.get(`${id} ${streamId} ${channel}`)]),
)
}
// ── Fixtures ───────────────────────────────────────────────────────────────
const addUser = (id, over = {}) => store.users.set(id, { id, role: 'player', status: 'active', ...over })
const optIn = (userId, streamId, channel, mode = 'instant') =>
store.prefs.set(`${userId} ${streamId} ${channel}`, mode)
let nextRuleId = 1
function addRule(over = {}) {
const id = nextRuleId++
const rule = {
id,
trigger_id: 'uo.house.idoc_warning',
name: `rule ${id}`,
enabled: true,
audience: 'owner',
audience_segment_id: null,
max_sends_per_hour: 100,
channels: ['email'],
template_keys: {},
conditions: null,
cooldown_seconds: 0,
delay_seconds: 0,
cancel_on: [],
...over,
}
store.rules.set(id, rule)
return rule
}
/** A validated event envelope, the shape `engagementEmit.emit` builds. */
const event = (over = {}) => ({
triggerId: 'uo.house.idoc_warning',
owner: 'uo',
version: 1,
subject: 'house-4001',
ownerUserId: 10,
dedupeKey: null,
occurredAt: T0.toISOString(),
data: { house: 'The Silver Anvil', decayStatus: 'IDOC' },
...over,
})
/** Register a batch, the way the loader's second pass commits one. */
function register(owner, fn) {
const api = registries.stage(owner)
fn(api)
registries.apply(api.staged)
}
const IDOC_TRIGGER = {
id: 'uo.house.idoc_warning',
label: 'House approaching collapse',
ceiling: 'owner',
audience: 'owner',
subjectKey: 'house',
variables: [
{ name: 'house', type: 'string', required: true, example: 'The Silver Anvil' },
{ name: 'decayStatus', type: 'string', required: false, example: 'IDOC' },
],
}
function registerUoTrigger(over = {}) {
register('uo', (api) => api.registerEventTriggers([{ ...IDOC_TRIGGER, ...over }]))
}
const outboxRows = (filter = () => true) => [...store.outbox.values()].filter(filter)
const scheduled = () => outboxRows((r) => r.status === 'scheduled')
// Core's channels register through `coreChannels`, the way app.js does.
// Requiring `channels` alone gets the empty map — that is the design, and the
// engine dropping every rule because no channel is registered is what a
// boot-order regression would look like.
function registerChannels() {
channels._reset()
delete require.cache[require.resolve('../src/engagement/coreChannels')]
// eslint-disable-next-line global-require
require('../src/engagement/coreChannels')
}
beforeEach(() => {
registries._reset()
registerChannels()
installStubs()
nextRuleId = 1
registerUoTrigger()
addUser(10)
// Every channel is opt-IN (§7.1 Q1, and channels.js `defaultMode: 'off'`), so
// a fixture that wants mail to happen has to say so. The test below that turns
// this off again is the one asserting exactly that.
optIn(10, 'uo.house.idoc_warning', 'email')
})
afterEach(() => {
registries._reset()
restoreOriginals()
})
// ── Acceptance: cooldowns ──────────────────────────────────────────────────
test('a trigger fired twice inside cooldown_seconds for the same (rule, user, subject) sends once', async () => {
addRule({ cooldown_seconds: 3600 })
await engine.dispatch(event(), T0)
await engine.dispatch(event(), later(60_000))
assert.equal(outboxRows().length, 1)
})
test('the same trigger for a DIFFERENT subject sends again — the multi-house case (§4.1)', async () => {
// The rule §4.1 warns about is "one IDOC mail per player per day": a player
// with four houses decaying should hear about all four, once each. Cooling on
// (rule, user) alone silently drops three of them, and this is the test that
// would fail if `subject_key` were ever dropped from the primary key.
addRule({ cooldown_seconds: 86_400 })
await engine.dispatch(event({ subject: 'house-4001' }), T0)
await engine.dispatch(event({ subject: 'house-4002' }), later(1000))
await engine.dispatch(event({ subject: 'house-4003' }), later(2000))
// …and the first house again, still inside the day.
await engine.dispatch(event({ subject: 'house-4001' }), later(3000))
assert.deepEqual(outboxRows().map((r) => r.subject_key).sort(), ['house-4001', 'house-4002', 'house-4003'])
})
test('a cooldown that has expired lets the same subject through again', async () => {
addRule({ cooldown_seconds: 60 })
await engine.dispatch(event(), T0)
await engine.dispatch(event(), later(61_000))
assert.equal(outboxRows().length, 2)
})
test('two rules on one trigger each get their own cooldown', async () => {
addRule({ cooldown_seconds: 3600 })
addRule({ cooldown_seconds: 3600 })
await engine.dispatch(event(), T0)
assert.equal(outboxRows().length, 2)
})
// ── Acceptance: dedupe ─────────────────────────────────────────────────────
test('a duplicate dedupe_key is a successful no-op, not a second row and not an error', async () => {
addRule()
const first = await engine.dispatch(event({ dedupeKey: 'idoc:4001:2026-08-29' }), T0)
const replay = await engine.dispatch(event({ dedupeKey: 'idoc:4001:2026-08-29' }), later(1000))
assert.equal(first.enqueued, 1)
assert.equal(replay.enqueued, 0)
assert.equal(replay.deduped, 1)
assert.equal(outboxRows().length, 1)
})
test('one dedupe_key fans out to every recipient — the key is scoped, not global', async () => {
// §4.2a's `UNIQUE (dedupe_key)` was a defect: a dedupe key names the EVENT, and
// one event legitimately becomes one row per (rule, user, channel). A global
// unique index would have let the FIRST recipient's row in and silently dropped
// everyone else's, which is the opposite of what dedupe is for.
addUser(11)
addUser(12)
for (const id of [10, 11, 12]) {
optIn(id, 'uo.house.idoc_warning', 'email')
optIn(id, 'uo.house.idoc_warning', 'inapp')
}
registries._reset()
registerUoTrigger({ ceiling: 'subscribers', audience: 'subscribers' })
addRule({ audience: 'subscribers', channels: ['email', 'inapp'] })
const result = await engine.dispatch(event({ dedupeKey: 'idoc:4001' }), T0)
// three users x two channels
assert.equal(result.enqueued, 6)
assert.equal(new Set(outboxRows().map((r) => r.dedupe_key)).size, 1)
})
// ── Acceptance: the grace window and cancellation ──────────────────────────
test('a scheduled row is cancelled by a cancel_on trigger and never sends', async () => {
register('uo', (api) =>
api.registerEventTriggers([
{ ...IDOC_TRIGGER, id: 'uo.house.repaired', label: 'House repaired', subjectKey: 'house' },
]),
)
const rule = addRule({ delay_seconds: 1800, cancel_on: ['uo.house.repaired'] })
await engine.dispatch(event(), T0)
assert.equal(scheduled().length, 1)
const cancelled = await engine.dispatch(
event({ triggerId: 'uo.house.repaired', subject: 'house-4001' }),
later(60_000),
)
assert.equal(cancelled.cancelled, 1)
// The window has passed; the worker finds nothing to do.
await worker.tick(later(1_900_000))
assert.equal(store.outbox.get(1).status, 'cancelled')
assert.equal(store.sends.length, 0)
assert.equal(rule.id, 1)
})
test('a resolving event with no owner cancels every recipient queued about that subject', async () => {
register('uo', (api) =>
api.registerEventTriggers([
{ ...IDOC_TRIGGER, id: 'uo.house.repaired', label: 'House repaired', ceiling: 'authenticated', audience: 'authenticated' },
]),
)
addUser(11)
registries._reset()
registerUoTrigger({ ceiling: 'authenticated', audience: 'authenticated' })
register('uo', (api) =>
api.registerEventTriggers([
{ ...IDOC_TRIGGER, id: 'uo.house.repaired', label: 'House repaired', ceiling: 'authenticated', audience: 'authenticated' },
]),
)
optIn(10, 'uo.house.idoc_warning', 'email')
optIn(11, 'uo.house.idoc_warning', 'email')
addRule({ audience: 'authenticated', delay_seconds: 600, cancel_on: ['uo.house.repaired'] })
await engine.dispatch(event(), T0)
assert.equal(scheduled().length, 2)
await engine.dispatch(
event({ triggerId: 'uo.house.repaired', subject: 'house-4001', ownerUserId: null }),
later(1000),
)
assert.equal(scheduled().length, 0)
})
test('cancellation leaves an in-flight row alone', async () => {
register('uo', (api) =>
api.registerEventTriggers([{ ...IDOC_TRIGGER, id: 'uo.house.repaired', label: 'House repaired' }]),
)
addRule({ delay_seconds: 600, cancel_on: ['uo.house.repaired'] })
await engine.dispatch(event(), T0)
// A worker has claimed it: cancelling now would leave two writers on one row.
await outboxDb.claim(1)
const result = await engine.dispatch(event({ triggerId: 'uo.house.repaired' }), later(1000))
assert.equal(result.cancelled, 0)
assert.equal(store.outbox.get(1).status, 'sending')
})
// ── Acceptance: exactly once across a restart ──────────────────────────────
test('a restart mid-window still sends exactly once', async () => {
addRule({ delay_seconds: 600 })
await engine.dispatch(event(), T0)
// "Restart" is the engine losing its process between enqueue and due_at. The
// outbox is the durable half, so the only question is whether the sweep after
// the restart double-delivers — and the CAS claim is what says it cannot.
await worker.tick(later(500_000)) // not yet due
assert.equal(store.sends.length, 0)
await worker.tick(later(700_000))
await worker.tick(later(700_001)) // a second instance, or the next tick
assert.equal(store.sends.length, 1)
})
test('two sweepers racing one due row: exactly one claim wins', async () => {
addRule()
await engine.dispatch(event(), T0)
// Both see the same candidate — findDue does not claim — and then disagree
// harmlessly about which of them owns it. §7.1 Q2's answer, as a test.
const [a] = await outboxDb.findDue(later(1000))
const [b] = await outboxDb.findDue(later(1000))
assert.equal(a.id, b.id)
assert.equal(await outboxDb.claim(a.id), true)
assert.equal(await outboxDb.claim(b.id), false)
})
// ── The send log ───────────────────────────────────────────────────────────
test('a row whose channel has no deliver() finishes failed, and the send log says why', async () => {
// `inapp`, because as of Phase 6 `email` DOES deliver. The inbox arrives in
// Phase 7, and until then recording 'sent' would be a lie in the one table
// whose purpose is answering "did they get it".
addRule({ channels: ['inapp'] })
optIn(10, 'uo.house.idoc_warning', 'inapp')
await engine.dispatch(event(), T0)
await worker.tick(later(1000))
assert.equal(store.outbox.get(1).status, 'failed')
assert.equal(store.sends.length, 1)
assert.equal(store.sends[0].status, 'failed')
assert.match(store.sends[0].detail, /no delivery implementation/)
assert.equal(store.sends[0].user_id, 10)
})
test('a transient failure is retried, and then given up on', async () => {
addRule()
await engine.dispatch(event(), T0)
let attempts = 0
const failing = async () => {
attempts += 1
return { outcome: 'retry', detail: 'smtp timeout' }
}
let at = later(1000)
for (let i = 0; i < worker.MAX_ATTEMPTS + 2; i += 1) {
// eslint-disable-next-line no-await-in-loop
const [row] = await outboxDb.findDue(at)
if (!row) break
// eslint-disable-next-line no-await-in-loop
await worker.processRow(row, at, failing)
at = new Date(at.getTime() + worker.RETRY_MS + 1000)
}
// Tried MAX_ATTEMPTS times and then stopped, rather than retrying forever.
assert.equal(attempts, worker.MAX_ATTEMPTS)
assert.equal(store.outbox.get(1).status, 'failed')
assert.equal(store.sends.length, 1)
assert.match(store.sends[0].detail, /smtp timeout/)
})
// ── Preferences ────────────────────────────────────────────────────────────
test("a user whose mode is 'off' for the channel is not enqueued", async () => {
registries._reset()
registerUoTrigger({ ceiling: 'authenticated', audience: 'authenticated' })
addUser(11)
optIn(10, 'uo.house.idoc_warning', 'email', 'instant')
optIn(11, 'uo.house.idoc_warning', 'email', 'off')
addRule({ audience: 'authenticated' })
await engine.dispatch(event(), T0)
assert.deepEqual(outboxRows().map((r) => r.user_id), [10])
})
test('absence means the CHANNEL default, and all three of core default off', async () => {
registries._reset()
registerUoTrigger({ ceiling: 'authenticated', audience: 'authenticated' })
store.prefs.clear()
addRule({ audience: 'authenticated' })
// Nobody has expressed anything, and email defaults 'off' (§3.1) — so an
// `authenticated` rule reaches nobody until people opt in. That is opt-IN
// working, not the engine failing.
const result = await engine.dispatch(event(), T0)
assert.equal(result.enqueued, 0)
assert.equal(channels.defaultMode('email'), 'off')
})
// **This reverses what Phase 4a asserted here**, and the reversal is Phase 6's
// §4.2b decision rather than a change of mind about queues. A digest is
// re-derived from the source tables at send time — that is what makes a hidden
// post absent from it and a user who lost access unreachable by it — so an outbox
// row for a digest recipient would be a second copy of the content with none of
// those properties. Nothing drains it, so nothing writes it.
test("a 'digest' preference does NOT enqueue — the digest re-derives at send time", async () => {
registries._reset()
registerUoTrigger({ ceiling: 'authenticated', audience: 'authenticated' })
optIn(10, 'uo.house.idoc_warning', 'email', 'digest')
addRule({ audience: 'authenticated' })
await engine.dispatch(event(), T0)
assert.equal(outboxRows().length, 0)
})
// ── The hourly ceiling (§7.1 Q3) ───────────────────────────────────────────
test('a rule stops at its hourly send ceiling', async () => {
registries._reset()
registerUoTrigger({ ceiling: 'authenticated', audience: 'authenticated' })
for (let i = 20; i < 30; i += 1) {
addUser(i)
optIn(i, 'uo.house.idoc_warning', 'email')
}
const rule = addRule({ audience: 'authenticated', max_sends_per_hour: 4 })
const result = await engine.dispatch(event(), T0)
// Eleven eligible recipients (the ten here plus the fixture's user 10), and a
// ceiling of four: four rows, and the rest are counted and dropped rather than
// queued for later - a rule at its ceiling is a rule an operator has to fix.
assert.equal(result.enqueued, 4)
assert.equal(result.enqueued + result.capped, 11)
assert.equal(rule.max_sends_per_hour, 4)
})
test('the hourly ceiling counts sends, not attempts', async () => {
// A broken transport must not silently consume a rule's whole budget and mute
// it: only rows the log records as 'sent' count against the ceiling.
const rule = addRule({ max_sends_per_hour: 2 })
store.sends.push({ rule_id: rule.id, status: 'failed', created_at: T0 })
store.sends.push({ rule_id: rule.id, status: 'suppressed', created_at: T0 })
const result = await engine.dispatch(event(), later(1000))
assert.equal(result.enqueued, 1)
})
// ── Ceilings: the security boundary, both halves ───────────────────────────
test('a rule may not be SAVED with an audience wider than its trigger permits', async () => {
const checked = await rules.validate({
triggerId: 'uo.house.idoc_warning', // ceiling: owner
name: 'IDOC warning',
channels: ['email'],
audience: 'authenticated',
})
assert.equal(checked.ok, false)
assert.match(checked.errors.join(' '), /wider than trigger/)
})
test('the ceiling is re-checked at SEND time, so a module narrowing its declaration stops a saved rule', async () => {
// The only way this can fail is the case it exists for: the rule was saved
// when the trigger permitted `authenticated`, and a module upgrade has since
// narrowed the declaration to `owner`. A save-time check alone would keep
// mailing the wider set forever.
registries._reset()
registerUoTrigger({ ceiling: 'authenticated', audience: 'authenticated' })
addUser(11)
optIn(10, 'uo.house.idoc_warning', 'email')
optIn(11, 'uo.house.idoc_warning', 'email')
addRule({ audience: 'authenticated' })
const before = await engine.dispatch(event(), T0)
assert.equal(before.enqueued, 2)
registries._reset()
registerUoTrigger({ ceiling: 'owner', audience: 'owner' }) // the upgrade
const after2 = await engine.dispatch(event({ subject: 'house-9' }), later(1000))
assert.equal(after2.enqueued, 0)
})
test('a rule for an unregistered trigger is dormant, not deleted and not an error', async () => {
const rule = addRule({ trigger_id: 'uo.gone.away' })
const listed = await rules.listAnnotated()
const found = listed.find((r) => r.id === rule.id)
assert.equal(found.dormant, true)
assert.match(found.dormantReasons.join(' '), /not registered/)
})
// ── Segments (§5.1a) ───────────────────────────────────────────────────────
function registerAudiences() {
register('uo', (api) =>
api.registerAudiences([
{ id: 'uo.team.members', label: 'Team members', ceiling: 'members', params: [{ id: 'teamId', type: 'int', required: true }], resolve: async ({ teamId }) => (teamId === 1 ? [10, 11] : [12]) },
{ id: 'uo.governors', label: 'Governors', ceiling: 'members', resolve: async () => [11, 12] },
{ id: 'uo.watchers', label: 'Watchers', ceiling: 'authenticated', resolve: async () => [10, 13] },
{ id: 'uo.flagged', label: 'Flagged accounts', ceiling: 'staff', resolve: async () => [10] },
]),
)
}
test('OR takes the TIGHTER ceiling — union-widens is the wrong implementation', async () => {
registerAudiences()
const checked = segments.validate({
op: 'or',
nodes: [{ audienceId: 'uo.governors' }, { audienceId: 'uo.watchers' }],
})
assert.equal(checked.ok, true)
// members is below authenticated, so the meet is members — NOT authenticated,
// which is what a "widest wins" reading would have given.
assert.equal(checked.ceiling, 'members')
})
test('two incomparable ceilings are refused rather than resolved to a guess', async () => {
registerAudiences()
const checked = segments.validate({
op: 'and',
nodes: [{ audienceId: 'uo.governors' }, { audienceId: 'uo.flagged' }],
})
assert.equal(checked.ok, false)
assert.match(checked.errors.join(' '), /no common ceiling/)
})
test('NOT does not constrain the ceiling — excluding people cannot widen', async () => {
registerAudiences()
// `members AND NOT staff` reaches strictly fewer people than `members`. If the
// complement's ceiling were folded into the meet, meet('members','staff') is
// null and this safe segment would be refused.
const checked = segments.validate({
op: 'and',
nodes: [{ audienceId: 'uo.governors' }, { op: 'not', nodes: [{ audienceId: 'uo.flagged' }] }],
})
assert.equal(checked.ok, true)
assert.equal(checked.ceiling, 'members')
})
test('NOT outside an AND is refused — a complement needs a set to take it from', async () => {
registerAudiences()
for (const expression of [
{ op: 'not', nodes: [{ audienceId: 'uo.governors' }] },
{ op: 'or', nodes: [{ audienceId: 'uo.governors' }, { op: 'not', nodes: [{ audienceId: 'uo.flagged' }] }] },
]) {
const checked = segments.validate(expression)
assert.equal(checked.ok, false)
assert.match(checked.errors.join(' '), /only allowed inside an "and"/)
}
})
test('a segment resolves through the module resolvers, and AND NOT subtracts', async () => {
registerAudiences()
const checked = segments.validate({
op: 'and',
nodes: [
{ audienceId: 'uo.team.members', params: { teamId: 1 } }, // [10, 11]
{ op: 'not', nodes: [{ audienceId: 'uo.flagged' }] }, // [10]
],
})
const resolved = await segments.resolve(checked.expression)
assert.equal(resolved.dormant, false)
assert.deepEqual(resolved.userIds, [11])
})
test('a segment whose module is uninstalled is DORMANT and sends to nobody', async () => {
registerAudiences()
const checked = segments.validate({ op: 'or', nodes: [{ audienceId: 'uo.governors' }, { audienceId: 'uo.watchers' }] })
store.segments.set(1, { id: 1, name: 'staff-ish', expression: checked.expression, ceiling: 'members' })
addUser(11)
optIn(10, 'uo.house.idoc_warning', 'email')
optIn(11, 'uo.house.idoc_warning', 'email')
registries._reset()
registerUoTrigger({ ceiling: 'members', audience: 'members' })
addRule({ audience: 'members', audience_segment_id: 1 })
// The module is gone: `resolveAudience` answers dormant + empty, and the rule
// must NOT fall back to anything. Reaching a different population than the one
// composed is the failure §5.1a rule 4 forbids.
const result = await engine.dispatch(event(), T0)
assert.equal(result.enqueued, 0)
assert.equal(outboxRows().length, 0)
})
test('a rule pointing at a deleted segment is dormant, never a fallback to its plain audience', async () => {
registries._reset()
registerUoTrigger({ ceiling: 'authenticated', audience: 'authenticated' })
optIn(10, 'uo.house.idoc_warning', 'email')
addRule({ audience: 'authenticated', audience_segment_id: 99 }) // no such segment
const result = await engine.dispatch(event(), T0)
assert.equal(result.enqueued, 0)
})
test("a plain 'members' audience with no segment reaches nobody", async () => {
registries._reset()
registerUoTrigger({ ceiling: 'members', audience: 'members' })
optIn(10, 'uo.house.idoc_warning', 'email')
addRule({ audience: 'members' })
const result = await engine.dispatch(event(), T0)
assert.equal(result.enqueued, 0)
})
// ── Conditions ─────────────────────────────────────────────────────────────
test('a condition narrows which firings are interesting', async () => {
addRule({
conditions: { variable: 'decayStatus', cmp: 'in', value: ['Greatly damaged', 'IDOC'] },
})
await engine.dispatch(event({ data: { house: 'A', decayStatus: 'IDOC' } }), T0)
await engine.dispatch(event({ subject: 'house-2', data: { house: 'B', decayStatus: 'LikeNew' } }), later(1000))
assert.equal(outboxRows().length, 1)
})
test('a condition naming a variable the trigger does not declare is refused at save, with the name', async () => {
const checked = await rules.validate({
triggerId: 'uo.house.idoc_warning',
name: 'typo',
channels: ['email'],
audience: 'owner',
conditions: { variable: 'decaystatus', cmp: 'eq', value: 'IDOC' },
})
assert.equal(checked.ok, false)
assert.match(checked.errors.join(' '), /"decaystatus" is not a variable/)
})
test('an absent variable makes every comparison false — including "is not"', async () => {
// `ne` is the one that tempts otherwise: "not equal to IDOC" reads as satisfied
// by nothing at all, and treating it that way would fire the rule on every
// event that omits an optional variable.
const c = { variable: 'decayStatus', cmp: 'ne', value: 'IDOC' }
assert.equal(conditions.evaluate(c, { house: 'A' }), false)
assert.equal(conditions.evaluate(c, { house: 'A', decayStatus: 'LikeNew' }), true)
assert.equal(conditions.evaluate({ variable: 'decayStatus', cmp: 'absent' }, { house: 'A' }), true)
})
test('a condition tree that no longer parses fails CLOSED', async () => {
// A stored condition that stops making sense must stop the mail, not decay
// into "no conditions" and reach everyone the rule could ever reach.
assert.equal(conditions.evaluate({ op: 'xor', nodes: [] }, {}), false)
assert.equal(conditions.evaluate('nonsense', {}), false)
assert.equal(conditions.evaluate(null, {}), true)
})
test('and / or / not compose', () => {
const data = { house: 'The Silver Anvil', decayStatus: 'IDOC' }
assert.equal(
conditions.evaluate(
{ op: 'and', nodes: [{ variable: 'decayStatus', cmp: 'eq', value: 'IDOC' }, { variable: 'house', cmp: 'contains', value: 'Silver' }] },
data,
),
true,
)
assert.equal(
conditions.evaluate({ op: 'not', nodes: [{ variable: 'decayStatus', cmp: 'eq', value: 'IDOC' }] }, data),
false,
)
assert.equal(
conditions.evaluate(
{ op: 'or', nodes: [{ variable: 'decayStatus', cmp: 'eq', value: 'LikeNew' }, { variable: 'house', cmp: 'startsWith', value: 'The' }] },
data,
),
true,
)
})
test('an operator cannot be applied to a type it does not fit', () => {
const declaration = registries.eventTrigger('uo.house.idoc_warning')
const checked = conditions.validate(declaration, { variable: 'house', cmp: 'gt', value: 'x' })
assert.equal(checked.ok, false)
assert.match(checked.errors.join(' '), /cannot be applied to a string/)
})
// ── Rule validation, the rest ──────────────────────────────────────────────
test('a new rule is created disabled unless it says otherwise (§7.1 Q3)', async () => {
const checked = await rules.validate({
triggerId: 'uo.house.idoc_warning',
name: 'IDOC warning',
channels: ['email'],
audience: 'owner',
})
assert.equal(checked.ok, true)
assert.equal(checked.rule.enabled, false)
assert.equal(checked.rule.max_sends_per_hour, 100)
})
test('a rule naming an unregistered channel is refused', async () => {
const checked = await rules.validate({
triggerId: 'uo.house.idoc_warning',
name: 'IDOC warning',
channels: ['carrier-pigeon'],
audience: 'owner',
})
assert.equal(checked.ok, false)
assert.match(checked.errors.join(' '), /no channel "carrier-pigeon"/)
})
test('the hourly ceiling has a hard upper bound an operator cannot type past', async () => {
const checked = await rules.validate({
triggerId: 'uo.house.idoc_warning',
name: 'IDOC warning',
channels: ['email'],
audience: 'owner',
maxSendsPerHour: 10_000_000,
})
assert.equal(checked.ok, false)
assert.match(checked.errors.join(' '), /maxSendsPerHour/)
})
test('cancelOn without a delay is refused — there is no window to cancel in', async () => {
const checked = await rules.validate({
triggerId: 'uo.house.idoc_warning',
name: 'IDOC warning',
channels: ['email'],
audience: 'owner',
cancelOn: ['uo.house.repaired'],
})
assert.equal(checked.ok, false)
assert.match(checked.errors.join(' '), /no effect without a delaySeconds/)
})
test('a rule whose channel was removed is dormant but still editable', async () => {
const rule = addRule({ channels: ['email', 'carrier-pigeon'] })
const listed = await rules.listAnnotated()
assert.equal(listed.find((r) => r.id === rule.id).dormant, true)
// …and only the live channel is used when it fires.
optIn(10, 'uo.house.idoc_warning', 'email')
await engine.dispatch(event(), T0)
assert.deepEqual([...new Set(outboxRows().map((r) => r.channel))], ['email'])
})
// ── The dispatch contract ──────────────────────────────────────────────────
test('dispatch never throws at its caller, even when the database is gone', async () => {
addRule()
rulesDb.enabledForTrigger = async () => {
throw new Error('connection lost')
}
const result = await engine.dispatch(event(), T0)
assert.equal(result.enqueued, 0)
})
test('an event nobody has written a rule for is a no-op', async () => {
const result = await engine.dispatch(event(), T0)
assert.equal(result.rules, 0)
assert.equal(outboxRows().length, 0)
})
test('a disabled rule does not fire', async () => {
addRule({ enabled: false })
const result = await engine.dispatch(event(), T0)
assert.equal(result.rules, 0)
})