File size: 12,662 Bytes
ee888e1
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
// Pure helpers for the digest cron's per-user compose loop.
//
// Extracted from scripts/seed-digest-notifications.mjs so they can be
// unit-tested without dragging the cron's env-checking side effects
// (DIGEST_CRON_ENABLED check, Upstash REST helper, Convex relay
// auth) into the test runtime. The cron imports back from here.

import { compareRules, MAX_STORIES_PER_USER } from './brief-compose.mjs';
import { generateDigestProse } from './brief-llm.mjs';

/**
 * Build the email subject string. Extracted so the synthesis-level
 * β†’ subject ternary can be unit-tested without standing up the whole
 * cron loop. (Plan acceptance criterion A6.i.)
 *
 * Rules:
 *   - synthesisLevel 1 or 2 + non-empty briefLead β†’ "Intelligence Brief"
 *   - synthesisLevel 3 OR empty/null briefLead β†’ "Digest"
 *
 * Mirrors today's UX where the editorial subject only appeared when
 * a real LLM-produced lead was available; the L3 stub falls back to
 * the plain "Digest" subject to set reader expectations correctly.
 *
 * @param {{ briefLead: string | null | undefined; synthesisLevel: number; shortDate: string }} input
 * @returns {string}
 */
export function subjectForBrief({ briefLead, synthesisLevel, shortDate }) {
  if (briefLead && synthesisLevel >= 1 && synthesisLevel <= 2) {
    return `WorldMonitor Intelligence Brief β€” ${shortDate}`;
  }
  return `WorldMonitor Digest β€” ${shortDate}`;
}

/**
 * Single source of truth for the digest's story window. Used by BOTH
 * the compose path (digestFor closure in the cron) and the send loop.
 * Without this, the brief lead can be synthesized from a 24h pool
 * while the channel body ships 7d / 12h of stories β€” reintroducing
 * the cross-surface divergence the canonical-brain refactor is meant
 * to eliminate, just in a different shape.
 *
 * `lastSentAt` is the rule's previous successful send timestamp (ms
 * since epoch) or null on first send. `defaultLookbackMs` is the
 * first-send fallback (today: 24h).
 *
 * @param {number | null | undefined} lastSentAt
 * @param {number} nowMs
 * @param {number} defaultLookbackMs
 * @returns {number}
 */
export function digestWindowStartMs(lastSentAt, nowMs, defaultLookbackMs) {
  return lastSentAt ?? (nowMs - defaultLookbackMs);
}

/**
 * Walk an annotated rule list and return the winning candidate +
 * its non-empty story pool. Two-pass: due rules first (so the
 * synthesis comes from a rule that's actually sending), then ALL
 * eligible rules (compose-only tick β€” keeps the dashboard brief
 * fresh for weekly/twice_daily users). Within each pass, walk by
 * compareRules priority and pick the FIRST candidate whose pool is
 * non-empty AND survives `tryCompose` (when provided).
 *
 * Returns null when every candidate is rejected β€” caller skips the
 * user (same as today's behavior on empty-pool exhaustion).
 *
 * Plan acceptance criteria A6.l (compose-only tick still works for
 * weekly user) + A6.m (winner walks past empty-pool top-priority
 * candidate). Codex Round-3 High #1 + Round-4 High #1 + Round-4
 * Medium #2.
 *
 * `tryCompose` (optional): called with `(cand, stories)` after a
 * non-empty pool is found. Returning a truthy value claims the
 * candidate as winner and the value is forwarded as `composeResult`.
 * Returning a falsy value (e.g. composeBriefFromDigestStories
 * dropped every story via its URL/headline/shape filters) walks to
 * the next candidate. Without this callback, the helper preserves
 * the original "first non-empty pool wins" semantics, which let a
 * filter-rejected top-priority candidate suppress the brief for the
 * user even when a lower-priority candidate would have shipped one.
 *
 * `digestFor` receives the full annotated candidate (not just the
 * rule) so callers can derive a per-candidate story window from
 * `cand.lastSentAt` β€” see `digestWindowStartMs`.
 *
 * `log` is the per-rejected-candidate log emitter β€” passed in so
 * tests can capture lines without reaching for console.log.
 *
 * @param {Array<{ rule: object; lastSentAt: number | null; due: boolean }>} annotated
 * @param {(cand: { rule: object; lastSentAt: number | null; due: boolean }) => Promise<unknown[] | null | undefined>} digestFor
 * @param {(line: string) => void} log
 * @param {string} userId
 * @param {((cand: { rule: object; lastSentAt: number | null; due: boolean }, stories: unknown[]) => Promise<unknown> | unknown)} [tryCompose]
 * @returns {Promise<{ winner: { rule: object; lastSentAt: number | null; due: boolean }; stories: unknown[]; composeResult?: unknown } | null>}
 */
export async function pickWinningCandidateWithPool(annotated, digestFor, log, userId, tryCompose) {
  if (!Array.isArray(annotated) || annotated.length === 0) return null;
  const sortedDue = annotated.filter((a) => a.due).sort((a, b) => compareRules(a.rule, b.rule));
  const sortedAll = [...annotated].sort((a, b) => compareRules(a.rule, b.rule));
  // Build the walk order, deduping by rule reference so the same
  // rule isn't tried twice (a due rule appears in both sortedDue
  // and sortedAll).
  const seen = new Set();
  const walkOrder = [];
  for (const cand of [...sortedDue, ...sortedAll]) {
    if (seen.has(cand.rule)) continue;
    seen.add(cand.rule);
    walkOrder.push(cand);
  }
  for (const cand of walkOrder) {
    const stories = await digestFor(cand);
    if (!stories || stories.length === 0) {
      log(
        `[digest] brief filter drops user=${userId} ` +
          `sensitivity=${cand.rule.sensitivity ?? 'high'} ` +
          `variant=${cand.rule.variant ?? 'full'} ` +
          `due=${cand.due} ` +
          `outcome=empty-pool ` +
          `in=0 dropped_severity=0 dropped_url=0 dropped_headline=0 dropped_shape=0 dropped_cap=0 out=0`,
      );
      continue;
    }
    if (typeof tryCompose === 'function') {
      const composeResult = await tryCompose(cand, stories);
      if (!composeResult) {
        log(
          `[digest] brief filter drops user=${userId} ` +
            `sensitivity=${cand.rule.sensitivity ?? 'high'} ` +
            `variant=${cand.rule.variant ?? 'full'} ` +
            `due=${cand.due} ` +
            `outcome=filter-rejected ` +
            `in=${stories.length} out=0`,
        );
        continue;
      }
      return { winner: cand, stories, composeResult };
    }
    return { winner: cand, stories };
  }
  return null;
}

/**
 * Run the three-level canonical synthesis fallback chain.
 *   L1: full pre-cap pool + ctx (profile, greeting, !public) β€” canonical.
 *   L2: envelope-sized slice + empty ctx β€” degraded fallback (mirrors
 *       today's enrichBriefEnvelopeWithLLM behaviour).
 *   L3: null synthesis β€” caller composes from stub.
 *
 * Returns { synthesis, level } with `synthesis` matching
 * generateDigestProse's output shape (or null on L3) and `level`
 * one of {1, 2, 3}.
 *
 * Pure helper β€” no I/O beyond the deps.callLLM the inner functions
 * already perform. Errors at L1 propagate to L2; L2 errors propagate
 * to L3 (null/stub). `trace` callback fires per level transition so
 * callers can quantify failure-mode distribution in production logs.
 *
 * Plan acceptance criterion A6.h (3-level fallback triggers).
 *
 * @param {string} userId
 * @param {Array} stories β€” full pre-cap pool
 * @param {string} sensitivity
 * @param {{ profile: string | null; greeting: string | null }} ctx
 * @param {{ callLLM: Function; cacheGet: Function; cacheSet: Function }} deps
 * @param {(level: 1 | 2 | 3, kind: 'success' | 'fall' | 'throw', err?: unknown) => void} [trace]
 * @returns {Promise<{ synthesis: object | null; level: 1 | 2 | 3 }>}
 */
export async function runSynthesisWithFallback(userId, stories, sensitivity, ctx, deps, trace) {
  const noteTrace = typeof trace === 'function' ? trace : () => {};
  // L1 β€” canonical
  try {
    const l1 = await generateDigestProse(userId, stories, sensitivity, deps, {
      profile: ctx?.profile ?? null,
      greeting: ctx?.greeting ?? null,
      isPublic: false,
    });
    if (l1) {
      noteTrace(1, 'success');
      return { synthesis: l1, level: 1 };
    }
    noteTrace(1, 'fall');
  } catch (err) {
    noteTrace(1, 'throw', err);
  }
  // L2 β€” degraded fallback
  try {
    const cappedSlice = (Array.isArray(stories) ? stories : []).slice(0, MAX_STORIES_PER_USER);
    const l2 = await generateDigestProse(userId, cappedSlice, sensitivity, deps);
    if (l2) {
      noteTrace(2, 'success');
      return { synthesis: l2, level: 2 };
    }
    noteTrace(2, 'fall');
  } catch (err) {
    noteTrace(2, 'throw', err);
  }
  // L3 β€” stub
  noteTrace(3, 'success');
  return { synthesis: null, level: 3 };
}

/**
 * READ-time freshness predicate. Returns true if the story:track:v1 row
 * should be dropped because its source `publishedAt` is older than the
 * cutoff. Used by buildDigest to keep pre-deploy residue (whose ingest
 * gate has since been tightened) from shipping in briefs.
 *
 * Behaviour matrix:
 *   - publishedAt is a positive integer epoch-ms AND < cutoff β†’ drop (true).
 *   - publishedAt is a positive integer epoch-ms AND β‰₯ cutoff β†’ keep (false).
 *   - publishedAt is missing/unparseable/zero/negative β†’ keep (false).
 *
 * The "missing β†’ keep" branch is back-compat for legacy story:track:v1
 * rows written before publishedAt was persisted. Pre-deploy residue
 * with no publishedAt is NOT caught here β€” handle it via the audit
 * script's `--mode=residue` (one-shot eviction). Once that has run AND
 * β‰₯1 cron cycle has refreshed publishedAt on still-active rows, any
 * row reaching this predicate without publishedAt is anomalous, not
 * residue.
 *
 * See: skill ingest-gate-tightening-leaves-residue-in-read-path.
 *
 * @param {Record<string, string> | null | undefined} track
 * @param {number} ageCutoffMs β€” drop rows with publishedAt strictly less than this
 * @returns {boolean}
 */
export function shouldDropTrackByAge(track, ageCutoffMs) {
  const pubMs = Number.parseInt(track?.publishedAt ?? '', 10);
  if (!Number.isInteger(pubMs) || pubMs <= 0) return false;
  return pubMs < ageCutoffMs;
}

/**
 * Compute the READ-time freshness cutoff for a given digest window.
 * Cutoff is anchored to `windowStartMs` (from `digestWindowStartMs`)
 * minus a 24h buffer that accommodates sustained stories whose first
 * mention sits just before the window edge.
 *
 * Daily user (24h window) β†’ 48h-ago cutoff.
 * Weekly user (7d window) β†’ 8d-ago cutoff.
 *
 * @param {number} windowStartMs
 * @returns {number}
 */
export function readTimeAgeCutoffMs(windowStartMs) {
  const STALE_BUFFER_MS = 24 * 60 * 60 * 1000;
  return windowStartMs - STALE_BUFFER_MS;
}

/**
 * Sprint 1 / U2 β€” option (a) canonical-send mapping.
 *
 * Given the per-user winner record from the compose phase
 * (briefByUser.get(userId), shape: `{ envelope, magazineUrl, chosenVariant,
 * synthesisLevel }`) and the candidate rule list for that user, return
 * the SINGLE rule whose channel body the send loop should use for this
 * user-slot. Under option (a), the email body and the magazine URL
 * BOTH come from this rule's pool β€” no per-rule fan-out, no
 * winner-vs-non-winner channel divergence.
 *
 * Returns null when:
 *   - briefByUser has no entry for this user (compose found no
 *     non-empty pool β€” nothing to send).
 *   - briefByUser entry has no chosenVariant (defensive; treat as
 *     "no canonical winner identified" and skip the send).
 *   - candidate list contains no rule whose `variant` matches
 *     `chosenVariant` (rule was deleted between compose and send;
 *     skip rather than misroute to a sibling rule).
 *
 * The send loop calls this BEFORE the per-rule isDue check β€” non-
 * winner rules are dropped here, so the cron's `digest:last-sent:v1:
 * ${userId}:${variant}` write only ever happens for the winner-variant
 * key. Old non-winner-variant keys from pre-option-(a) sends remain in
 * Redis until their 8d TTL elapses; they are orphan but harmless
 * (nothing reads them after this change).
 *
 * Pure helper β€” no I/O.
 *
 * @param {{ chosenVariant?: string } | null | undefined} brief
 * @param {Array<{ userId?: string; variant?: string }>} userRules β€” all rules for ONE user
 * @returns {{ userId?: string; variant?: string } | null}
 */
export function selectCanonicalSendRule(brief, userRules) {
  if (!brief || typeof brief.chosenVariant !== 'string' || brief.chosenVariant === '') {
    return null;
  }
  if (!Array.isArray(userRules) || userRules.length === 0) return null;
  for (const rule of userRules) {
    if (!rule || typeof rule.variant !== 'string') continue;
    if (rule.variant === brief.chosenVariant) return rule;
  }
  return null;
}