From e215664b663490444cd983380fa0dc115855124e Mon Sep 17 00:00:00 2001 From: Adam Eivy Date: Thu, 3 Sep 2026 22:45:27 +0000 Subject: [PATCH] perf: memoize handle resolution and cache the unanswered-thread scan (#6025) The Tribe "unanswered thread" detector ran a full timeline scan on every call: the dashboard alerts widget polls it every 120s and the Tribe page hits it again on mount, so back-to-back callers each paid for up to eight 2,000-row queries against the activity timeline plus a handle-resolution pass over every inbound turn. - Resolve counterpart handles through one memoizing resolver per scan (`createHandleResolver`), so a busy 1:1 thread's repeated handle runs the handle regexes, phone/email normalization and the Tribe/Contacts match once instead of once per message. `enrichActivityEvent` takes the resolver as an optional argument and behaves identically without it. - Cache the detection pass for 120s, keyed by detection window and holding the unsliced thread list so callers asking for different limits share one result. Caching the promise also collapses concurrent callers onto a single scan; a failed pass is never cached. Detected threads, ordering and outreach-draft behavior are unchanged. Claude-Session: https://claude.ai/code/session_01RA3pD5YM2dukQwbZ3pC6WA --- server/services/identityResolve.js | 35 +++++++++- server/services/identityResolve.test.js | 21 ++++++ server/services/tribeOutreach.js | 46 ++++++++++++- server/services/tribeOutreach.test.js | 92 ++++++++++++++++++++----- 4 files changed, 172 insertions(+), 22 deletions(-) diff --git a/server/services/identityResolve.js b/server/services/identityResolve.js index 030d5e1386..36c60f910d 100644 --- a/server/services/identityResolve.js +++ b/server/services/identityResolve.js @@ -107,6 +107,29 @@ export async function loadResolverContext() { }; } +/** + * Memoizing `resolveHandle` bound to ONE context (#6025). + * + * `resolveHandle` is pure for a given ctx but not cheap — it runs the handle + * regexes, phone/email normalization, a Tribe index match and a Contacts lookup + * on every call. Bulk callers (the outreach scan, timeline enrichment) hit the + * same handful of handles thousands of times per pass, so hand them one resolver + * whose Map collapses that to once per distinct handle. + * + * The cache lives as long as the returned function — build one per pass so it + * can never outlive the ctx it resolves against. + */ +export function createHandleResolver(ctx) { + const cache = new Map(); + return (handle) => { + const key = handle == null ? '' : String(handle); + if (cache.has(key)) return cache.get(key); + const res = resolveHandle(key, ctx); + cache.set(key, res); + return res; + }; +} + /** * Resolve many handles with one shared context (avoids N Tribe loads). */ @@ -153,15 +176,21 @@ export function enrichConversationRow(row, ctx) { /** * Enrich activity events: attach counterpart displayName on participants + title. + * + * `resolve` optionally overrides how a handle is resolved — pass a + * `createHandleResolver(ctx)` when enriching many events so repeated handles are + * resolved once instead of once per event (#6025). Defaults to an uncached + * `resolveHandle` against `ctx`. */ -export function enrichActivityEvent(event, ctx) { +export function enrichActivityEvent(event, ctx, resolve) { if (!event) return event; + const resolveOne = resolve || ((h) => resolveHandle(h, ctx)); const handle = event.metadata?.handle || ''; - const res = resolveHandle(handle, ctx); + const res = resolveOne(handle); const participants = (event.participants || []).map((p) => { const key = p.phone || p.email || ''; if (!key) return p; - const pr = resolveHandle(key, ctx); + const pr = resolveOne(key); if (!pr.displayName) return p; return { ...p, diff --git a/server/services/identityResolve.test.js b/server/services/identityResolve.test.js index 25259c6ac9..a4014aee20 100644 --- a/server/services/identityResolve.test.js +++ b/server/services/identityResolve.test.js @@ -6,6 +6,7 @@ import { displayLabel, enrichConversationRow, enrichActivityEvent, + createHandleResolver, } from './identityResolve.js'; describe('identityResolve', () => { @@ -73,4 +74,24 @@ describe('identityResolve', () => { expect(ev.participants[0].name).toBe('Tribe Friend'); expect(ev.participants[0].personId).toBe('p1'); }); + + // #6025: bulk callers (the outreach scan) hand enrichActivityEvent one memoizing + // resolver so a repeated counterpart handle is matched once, not once per event. + it('memoizes repeated handles and enriches identically through an injected resolver', () => { + const resolver = createHandleResolver(ctx); + expect(resolver('+15551234567')).toBe(resolver('+15551234567')); + expect(resolver('+15559876543').displayName).toBe('Contact Only'); + + let injectedCalls = 0; + const counting = (h) => { injectedCalls += 1; return resolver(h); }; + const event = { + title: '+15551234567', + metadata: { handle: '+15551234567' }, + participants: [{ phone: '+15551234567' }], + }; + // The injected resolver is used for every lookup, and the enrichment is + // byte-identical to the uncached default path. + expect(enrichActivityEvent(event, ctx, counting)).toEqual(enrichActivityEvent(event, ctx)); + expect(injectedCalls).toBeGreaterThan(0); + }); }); diff --git a/server/services/tribeOutreach.js b/server/services/tribeOutreach.js index 431807af8a..988c503620 100644 --- a/server/services/tribeOutreach.js +++ b/server/services/tribeOutreach.js @@ -261,17 +261,52 @@ export function groupUnansweredThreads(events, { now = Date.now(), staleAfterMs return out; } +// Detection-pass cache (#6025). The scan reads up to EVENT_CAP rows per scope +// across the timeline's largest table, and its verdict changes only when new +// messages land — but the dashboard alert widget polls every 120s and the Tribe +// page hits the same detection on mount, so back-to-back callers used to each pay +// for a full re-scan. One entry (keyed by detection window, holding the UNSLICED +// thread list so any `limit` is served from it) is enough: the callers differ only +// in `limit`. Caching the promise also collapses concurrent callers onto one pass. +export const DETECTION_CACHE_TTL_MS = 120000; +let detectionCache = null; // { key, expiresAt, promise } + +/** Drop the cached detection pass (tests, and any caller that must see fresh state). */ +export function invalidateUnansweredThreadsCache() { + detectionCache = null; +} + /** * Detect unanswered inbound threads from Tribe people via the activity timeline. * No LLM — pure timeline + Tribe reads. Returns [] on any read failure (a nudge * source that quietly yields nothing beats one that throws into the alert sweep). + * + * Results are cached for DETECTION_CACHE_TTL_MS per detection window; a failed + * pass is never cached. */ export async function findUnansweredTribeThreads({ withinDays = DEFAULT_WITHIN_DAYS, staleAfterHours = DEFAULT_STALE_HOURS, limit = DEFAULT_LIMIT, } = {}) { - const [{ listEvents }, { loadResolverContext, enrichActivityEvent }, { listAccounts }] = await Promise.all([ + const key = `${withinDays}:${staleAfterHours}`; + const now = Date.now(); + if (!detectionCache || detectionCache.key !== key || detectionCache.expiresAt <= now) { + const entry = { key, expiresAt: now + DETECTION_CACHE_TTL_MS, promise: null }; + entry.promise = detectUnansweredTribeThreads({ withinDays, staleAfterHours }) + .catch((err) => { + if (detectionCache === entry) detectionCache = null; + throw err; + }); + detectionCache = entry; + } + const threads = await detectionCache.promise; + return threads.slice(0, Math.max(0, limit)); +} + +/** One full detection pass — the uncached, unsliced scan behind the cache above. */ +async function detectUnansweredTribeThreads({ withinDays, staleAfterHours }) { + const [{ listEvents }, { loadResolverContext, enrichActivityEvent, createHandleResolver }, { listAccounts }] = await Promise.all([ import('./humanActivity.js'), import('./identityResolve.js'), import('./messageAccounts.js'), @@ -345,6 +380,10 @@ export async function findUnansweredTribeThreads({ // invisible, so its inbound must not surface (and its sent turns must not vouch // for another account's thread). const tagged = sent.filter((ev) => !ev.metadata?.isReaction && isTwoWay(ev)); + // One memoized resolver for the whole pass (#6025): a busy 1:1 thread repeats + // the same counterpart handle across every one of its inbound turns, and the + // uncached path re-ran the handle regexes + Tribe/Contacts match for each. + const resolveHandleOnce = createHandleResolver(ctx); for (const ev of received) { if (ev.metadata?.isReaction) continue; if (!isTwoWay(ev)) continue; @@ -356,7 +395,7 @@ export async function findUnansweredTribeThreads({ // Resolve the SENDER to a Tribe person via the event's counterpart handle // ONLY (a 1:1 received event's `metadata.handle` IS the sender, which // `enrichActivityEvent` resolves). An unresolved sender is skipped, not guessed. - const enriched = enrichActivityEvent(ev, ctx); + const enriched = enrichActivityEvent(ev, ctx, resolveHandleOnce); const personId = enriched.personId || null; if (!personId) continue; const ring = ringById.get(personId); @@ -369,7 +408,8 @@ export async function findUnansweredTribeThreads({ staleAfterMs: staleAfterHours * HOUR_MS, withinMs, }); - return threads.slice(0, Math.max(0, limit)); + console.log(`🤝 Outreach scan: ${threads.length} unanswered thread(s) from ${received.length} inbound turn(s) in ${Date.now() - now}ms`); + return threads; } // Map one timeline event to a synthetic message (buildThreadContext / diff --git a/server/services/tribeOutreach.test.js b/server/services/tribeOutreach.test.js index fc2fa543d4..3a5bd38dd2 100644 --- a/server/services/tribeOutreach.test.js +++ b/server/services/tribeOutreach.test.js @@ -1,6 +1,6 @@ import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest'; -import { groupUnansweredThreads, generateOutreachDraft, findUnansweredTribeThreads, buildTwoWayGate, outreachTemplateForSource } from './tribeOutreach.js'; +import { groupUnansweredThreads, generateOutreachDraft, findUnansweredTribeThreads, buildTwoWayGate, outreachTemplateForSource, invalidateUnansweredThreadsCache, DETECTION_CACHE_TTL_MS } from './tribeOutreach.js'; // generateOutreachDraft / findUnansweredTribeThreads dynamically import these — // mock them so the draft-side logic (idempotent reuse, anchoring the reply to the @@ -9,14 +9,14 @@ vi.mock('./tribe.js', () => ({ getPerson: vi.fn() })); vi.mock('./humanActivity.js', () => ({ listEvents: vi.fn() })); vi.mock('./messageEvaluator.js', () => ({ generateReplyBody: vi.fn() })); vi.mock('./messageDrafts.js', () => ({ createDraft: vi.fn(), listDrafts: vi.fn() })); -vi.mock('./identityResolve.js', () => ({ loadResolverContext: vi.fn(), enrichActivityEvent: vi.fn() })); +vi.mock('./identityResolve.js', () => ({ loadResolverContext: vi.fn(), enrichActivityEvent: vi.fn(), createHandleResolver: vi.fn() })); vi.mock('./messageAccounts.js', () => ({ listAccounts: vi.fn() })); import { getPerson } from './tribe.js'; import { listEvents } from './humanActivity.js'; import { generateReplyBody } from './messageEvaluator.js'; import { createDraft, listDrafts } from './messageDrafts.js'; -import { loadResolverContext, enrichActivityEvent } from './identityResolve.js'; +import { loadResolverContext, enrichActivityEvent, createHandleResolver } from './identityResolve.js'; import { listAccounts } from './messageAccounts.js'; // `groupUnansweredThreads` is the pure detection core — no DB, no LLM. These @@ -271,22 +271,28 @@ describe('outreachTemplateForSource (#2796)', () => { }); }); +// Shared detector setup for the findUnansweredTribeThreads suites below. +// findUnansweredTribeThreads() reads Date.now() directly (no injectable `now`), so +// pin the clock to NOW — otherwise the fixtures' ages drift with wall-clock time +// and the daysAgo(3) inbound falls outside the 14-day window once real time passes +// NOW+14d (this test began failing at 2026-07-29T12:00Z). Only Date is faked; the +// code under test uses no real timers. The detection cache (#6025) is module-level, +// so it must be dropped between tests or the second suite reads the first's result. +const primeDetector = () => { + vi.clearAllMocks(); + invalidateUnansweredThreadsCache(); + vi.useFakeTimers({ toFake: ['Date'] }); + vi.setSystemTime(NOW); + // One tribe person, resolved from the inbound's handle. + loadResolverContext.mockResolvedValue({ people: [{ id: 'p1', ring: 'tribe', name: 'Alex' }] }); + enrichActivityEvent.mockReturnValue({ personId: 'p1', displayName: 'Alex' }); + createHandleResolver.mockImplementation(() => vi.fn()); +}; + describe('findUnansweredTribeThreads — per-account email querying (#2820)', () => { const RECENT = new Date(NOW - 3600000).toISOString(); - beforeEach(() => { - vi.clearAllMocks(); - // findUnansweredTribeThreads() reads Date.now() directly (no injectable - // `now`), so pin the clock to NOW — otherwise the fixtures' ages drift with - // wall-clock time and the daysAgo(3) inbound falls outside the 14-day window - // once real time passes NOW+14d (this test began failing at 2026-07-29T12:00Z). - // Only Date is faked; the code under test uses no real timers. - vi.useFakeTimers({ toFake: ['Date'] }); - vi.setSystemTime(NOW); - // One tribe person, resolved from the inbound's handle. - loadResolverContext.mockResolvedValue({ people: [{ id: 'p1', ring: 'tribe', name: 'Alex' }] }); - enrichActivityEvent.mockReturnValue({ personId: 'p1', displayName: 'Alex' }); - }); + beforeEach(primeDetector); afterEach(() => { vi.useRealTimers(); @@ -337,6 +343,60 @@ describe('findUnansweredTribeThreads — per-account email querying (#2820)', () }); }); +describe('findUnansweredTribeThreads — detection cache + handle memoization (#6025)', () => { + beforeEach(primeDetector); + + afterEach(() => { + vi.useRealTimers(); + }); + + // Several unanswered 1:1 turns from the SAME counterpart handle, in distinct + // conversations (so each surfaces as its own thread). + const inboundFrom = (n) => ({ + kind: 'message.received', + source: 'imessage', + happenedAt: daysAgo(3), + summary: `ping ${n}`, + metadata: { chatGuid: `chat-${n}`, handle: '+15550001' }, + }); + const mockTimeline = (events) => { + listAccounts.mockResolvedValue([]); + listEvents.mockImplementation(async ({ source, kind }) => ( + source === 'imessage' && kind === 'message.received' ? events : [] + )); + }; + + it('shares ONE memoizing resolver across the pass so a repeated handle is matched once', async () => { + mockTimeline([inboundFrom(1), inboundFrom(2), inboundFrom(3)]); + + const threads = await findUnansweredTribeThreads(); + + expect(threads).toHaveLength(3); + expect(createHandleResolver).toHaveBeenCalledTimes(1); + const resolver = createHandleResolver.mock.results[0].value; + expect(enrichActivityEvent).toHaveBeenCalledTimes(3); + for (const call of enrichActivityEvent.mock.calls) expect(call[2]).toBe(resolver); + }); + + it('serves a repeat call from cache and re-scans only once the TTL lapses', async () => { + mockTimeline([inboundFrom(1), inboundFrom(2)]); + + const first = await findUnansweredTribeThreads(); + const queries = listEvents.mock.calls.length; + expect(queries).toBeGreaterThan(0); + + // The dashboard poll and the Tribe page ask for different limits — both are + // served from the one cached (unsliced) detection result. + const cached = await findUnansweredTribeThreads({ limit: 1 }); + expect(listEvents.mock.calls.length).toBe(queries); + expect(cached).toEqual(first.slice(0, 1)); + + vi.setSystemTime(NOW + DETECTION_CACHE_TTL_MS + 1); + await findUnansweredTribeThreads(); + expect(listEvents.mock.calls.length).toBeGreaterThan(queries); + }); +}); + describe('generateOutreachDraft', () => { beforeEach(() => { vi.clearAllMocks();