diff --git a/apps/web/package.json b/apps/web/package.json index ec05b5f..c1355ef 100644 --- a/apps/web/package.json +++ b/apps/web/package.json @@ -25,7 +25,8 @@ "env:setup-dev": "npx tsx src/lib/atproto/scripts/setup-dev.ts", "tunnel": "npx tsx src/lib/atproto/scripts/tunnel.ts", "publish-lexicons": "contrail-lex publish", - "seed:conference": "bun run scripts/publish-test-conference.ts" + "seed:conference": "bun run scripts/publish-test-conference.ts", + "geocode:backfill": "tsx scripts/geocode-events.ts" }, "devDependencies": { "@atcute/atproto": "^3.1.10", diff --git a/apps/web/scripts/geocode-cache.sql b/apps/web/scripts/geocode-cache.sql new file mode 100644 index 0000000..20c5f23 --- /dev/null +++ b/apps/web/scripts/geocode-cache.sql @@ -0,0 +1,14 @@ +-- Derived-coordinate cache for address-only events + the geocode job's worklist +-- and done-marker (Meili can't filter "missing _geo", so we track resolution +-- ourselves). Idempotent / restartable. Applied once to the openmeet-atmo D1; +-- the sink reads it, the external geocode job writes it. +CREATE TABLE IF NOT EXISTS geocode_cache ( + address_norm TEXT PRIMARY KEY, -- normalized address key (address-norm.ts) + lat REAL, -- NULL when unresolved (negative cache) + lng REAL, + precision TEXT, -- provider-reported granularity + source TEXT NOT NULL, -- e.g. 'locationiq' / 'nominatim' + geocoded_at INTEGER NOT NULL, -- epoch ms + fail_count INTEGER NOT NULL DEFAULT 0, + last_error TEXT +); diff --git a/apps/web/scripts/geocode-events.ts b/apps/web/scripts/geocode-events.ts new file mode 100644 index 0000000..b3413f2 --- /dev/null +++ b/apps/web/scripts/geocode-events.ts @@ -0,0 +1,263 @@ +// apps/web/scripts/geocode-events.ts +// External geocode job (runs OFF Cloudflare): finds address-only events in the +// openmeet-atmo D1 that lack coordinates, geocodes each unique address through +// the config-selected Nominatim/LocationIQ client, writes geocode_cache, and +// _geo-updates the affected Meili docs. ONE job: "find work" and "geocode" are +// sequential steps sharing the cache, not two passes. Rate limiting is the +// sleep between calls — the whole reason geocoding stays off Workers. +// +// Run (backfill MUST point the geocoder at LocationIQ, never public Nominatim): +// CLOUDFLARE_API_TOKEN=… MEILI_URL=https://search.testnet.openmeet.net MEILI_KEY=… \ +// GEOCODER_URL=https://us1.locationiq.com/v1/search GEOCODER_KEY=… \ +// pnpm -C apps/web exec tsx scripts/geocode-events.ts --limit 50 +import { createD1Client, type D1Client } from '../src/lib/search/server/d1-http'; +import { + createGeocoder, + addressToQuery, + requireGeocoderForBulk +} from '../src/lib/search/server/geocoder'; +import { + isEligible, + groupEventsByNorm, + addressNeedingGeocode, + type GeocodeCacheRow, + type WorklistEvent +} from '../src/lib/search/server/geocode-cache'; +import { eventToSearchDoc } from '../src/lib/search/server/normalize'; +import { discoverableSql } from '../src/lib/search/server/discoverability'; +import { MeiliEventIndex, EVENT_COLLECTION } from '../src/lib/search/server/meili-sink'; + +const env = process.env; +const argv = process.argv.slice(2); +const flag = (name: string) => argv.includes(name); +const opt = (name: string, def?: string) => { + const i = argv.indexOf(name); + return i >= 0 && i + 1 < argv.length ? argv[i + 1] : def; +}; + +// Strict: 0 = no cap, otherwise a positive integer. Reject negatives/non-integers +// up front — `Number('-1') || 0` is -1, which used to slip past BOTH the bulk +// guard (its old limit===0 check) and the cap (`limit > 0 ? slice : all`), +// silently running an uncapped keyless backfill against public Nominatim. +const parseLimit = (raw: string | undefined): number => { + const n = Number(raw); + if (!Number.isInteger(n) || n < 0) { + throw new Error( + `--limit must be a non-negative integer (0 = no cap); got ${JSON.stringify(raw)}` + ); + } + return n; +}; + +const retryNegative = flag('--retry-negative'); +const dryRun = flag('--dry-run'); +const allowPublicNominatim = flag('--allow-public-nominatim'); +const limit = parseLimit(opt('--limit', '0')); // 0 = no cap +const sleepMs = Number(env.GEOCODE_SLEEP_MS ?? '1100'); + +const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms)); + +// Worklist: discoverable events carrying a .address location. We deliberately +// do NOT exclude coordinate locations in SQL — that filtered by $type presence, +// which diverges from the sink's actual _geo derivation (it ignored fsq, and +// excluded events whose only geo/hthree coords are out of range and so never +// get an in-index _geo). The precise "does this already resolve to coordinates?" +// decision is made in memory by addressNeedingGeocode (recordGeo), the same +// derivation the sink uses. json_each walks locations[]; the "$type" key is +// quoted because it starts with $. +const WORKLIST_SQL = ` +SELECT r.uri AS uri, r.did AS did, r.rkey AS rkey, r.record AS record +FROM records_event AS r +WHERE EXISTS ( + SELECT 1 FROM json_each(r.record, '$.locations') + WHERE json_extract(value, '$."$type"') = 'community.lexicon.location.address') + AND ${discoverableSql('r.record')} +`; + +// Re-read events from records_event by uri, keeping only those still present and +// discoverable. The parsed record drives a full-doc upsert, so it's read fresh +// (not from the job-start worklist) — an event deleted or unlisted while the slow +// geocode loop runs is skipped here instead of being resurrected as a stale doc. +async function fetchLiveDocs( + d1: D1Client, + uris: string[] +): Promise<{ uri: string; did: string; rkey: string; record: Record }[]> { + if (uris.length === 0) return []; + const placeholders = uris.map(() => '?').join(','); + const rows = await d1.query<{ uri: string; did: string; rkey: string; record: string }>( + `SELECT uri, did, rkey, record FROM records_event + WHERE uri IN (${placeholders}) AND ${discoverableSql('record')}`, + uris + ); + const out: { uri: string; did: string; rkey: string; record: Record }[] = []; + for (const r of rows) { + try { + out.push({ uri: r.uri, did: r.did, rkey: r.rkey, record: JSON.parse(r.record) }); + } catch { + // Unparseable record JSON — skip, same as the worklist parse. + } + } + return out; +} + +async function main() { + if (!env.CLOUDFLARE_API_TOKEN) throw new Error('CLOUDFLARE_API_TOKEN is required'); + // No hardcoded account/DB fallbacks: the target D1 must be chosen explicitly so the + // job can never silently write a baked-in database, and so no infra IDs live in source. + if (!env.CLOUDFLARE_ACCOUNT_ID || !env.D1_DATABASE_ID) + throw new Error('CLOUDFLARE_ACCOUNT_ID and D1_DATABASE_ID are required'); + if (!env.MEILI_URL || !env.MEILI_KEY) throw new Error('MEILI_URL and MEILI_KEY are required'); + if (!env.GEOCODER_KEY) { + console.warn( + '[geocode] GEOCODER_KEY is unset → using PUBLIC Nominatim. OK for a small drip; ' + + 'a bulk backfill against public Nominatim risks a silent IP ban. Use LocationIQ for backfill.' + ); + } + // Hard-stop a bulk/uncapped keyless run before it can touch public Nominatim. + requireGeocoderForBulk({ + hasKey: !!env.GEOCODER_KEY, + dryRun, + limit, + allowPublic: allowPublicNominatim + }); + + const d1 = createD1Client({ + accountId: env.CLOUDFLARE_ACCOUNT_ID, + databaseId: env.D1_DATABASE_ID, + apiToken: env.CLOUDFLARE_API_TOKEN + }); + const geocoder = createGeocoder(env); + const meili = new MeiliEventIndex({ + url: env.MEILI_URL, + apiKey: env.MEILI_KEY, + indexUid: env.SEARCH_INDEX ?? 'events' + }); + + // Defensive: the table is created by geocode-cache.sql, but a fresh DB + // shouldn't make the job crash before it can self-heal. + await d1.query( + `CREATE TABLE IF NOT EXISTS geocode_cache ( + address_norm TEXT PRIMARY KEY, lat REAL, lng REAL, precision TEXT, + source TEXT NOT NULL, geocoded_at INTEGER NOT NULL, + fail_count INTEGER NOT NULL DEFAULT 0, last_error TEXT)` + ); + + // Load the whole cache once (small) and decide eligibility in memory. + const cacheRows = await d1.query(`SELECT * FROM geocode_cache`); + const cache = new Map(cacheRows.map((r) => [r.address_norm, r])); + + // Worklist → WorklistEvent[] (parse record JSON; keep only events that need + // geocoding — an address location AND no coordinates the index already + // derives, per addressNeedingGeocode). + const rawEvents = await d1.query<{ uri: string; did: string; rkey: string; record: string }>( + WORKLIST_SQL + ); + const events: WorklistEvent[] = []; + for (const e of rawEvents) { + let record: Record; + try { + record = JSON.parse(e.record); + } catch { + continue; + } + const loc = addressNeedingGeocode(record); + if (loc) events.push({ uri: e.uri, did: e.did, rkey: e.rkey, loc }); + } + + const byNorm = groupEventsByNorm(events); + const now = Date.now(); + const work = [...byNorm.entries()].filter(([norm]) => + isEligible(cache.get(norm), now, retryNegative) + ); + const capped = limit > 0 ? work.slice(0, limit) : work; + + console.log( + `[geocode] worklist events=${events.length} unique-addresses=${byNorm.size} ` + + `eligible=${work.length} processing=${capped.length}${dryRun ? ' (dry-run)' : ''}` + ); + + let resolved = 0; + let negative = 0; + let transient = 0; + let skippedGone = 0; + for (const [norm, group] of capped) { + const query = addressToQuery(group[0].loc); + if (dryRun) { + console.log(`[geocode] would geocode "${query}" -> ${group.length} event(s)`); + continue; + } + try { + const point = await geocoder.geocode(query); + if (point) { + await d1.query( + `INSERT INTO geocode_cache (address_norm, lat, lng, precision, source, geocoded_at, fail_count, last_error) + VALUES (?, ?, ?, ?, ?, ?, 0, NULL) + ON CONFLICT(address_norm) DO UPDATE SET + lat=excluded.lat, lng=excluded.lng, precision=excluded.precision, + source=excluded.source, geocoded_at=excluded.geocoded_at, fail_count=0, last_error=NULL`, + [ + norm, + point.lat, + point.lng, + point.precision ?? null, + env.GEOCODER_KEY ? 'locationiq' : 'nominatim', + now + ] + ); + // Re-read the events fresh and upsert the FULL doc with _geo attached, + // keeping only those still present AND discoverable. A full-doc upsert is + // idempotent and needs no "is it indexed?" snapshot: it merges if the + // event is already indexed, lands a complete doc if not (never a {id,_geo} + // stub), and is identical to what the sink writes — so a concurrent sink + // write converges instead of clobbering. + const live = await fetchLiveDocs( + d1, + group.map((e) => e.uri) + ); + skippedGone += group.length - live.length; + if (live.length) { + await meili.upsert( + live.map((r) => { + const doc = eventToSearchDoc({ + uri: r.uri, + did: r.did, + collection: EVENT_COLLECTION, + rkey: r.rkey, + record: r.record + }); + // Don't overwrite a coordinate _geo the record gained mid-run. + if (!doc._geo) doc._geo = { lat: point.lat, lng: point.lng }; + return doc; + }) + ); + } + resolved++; + } else { + // No-match: write/increment a negative row (backoff handled by isEligible). + await d1.query( + `INSERT INTO geocode_cache (address_norm, lat, lng, precision, source, geocoded_at, fail_count, last_error) + VALUES (?, NULL, NULL, NULL, ?, ?, 1, 'no match') + ON CONFLICT(address_norm) DO UPDATE SET + source=excluded.source, geocoded_at=excluded.geocoded_at, fail_count=geocode_cache.fail_count+1, last_error='no match'`, + [norm, env.GEOCODER_KEY ? 'locationiq' : 'nominatim', now] + ); + negative++; + } + } catch (err) { + // Transient (HTTP/network): DON'T write a negative row — retry next run. + transient++; + console.warn(`[geocode] transient error for "${query}": ${(err as Error).message}`); + } + await sleep(sleepMs); + } + + console.log( + `[geocode] done resolved=${resolved} negative=${negative} transient=${transient} ` + + `skipped-gone=${skippedGone}` + ); +} + +main().catch((e) => { + console.error('[geocode] fatal:', e); + process.exit(1); +}); diff --git a/apps/web/src/lib/contrail.config.ts b/apps/web/src/lib/contrail.config.ts index c54943d..5ef7011 100644 --- a/apps/web/src/lib/contrail.config.ts +++ b/apps/web/src/lib/contrail.config.ts @@ -2,6 +2,7 @@ import type { ContrailConfig } from '@atmo-dev/contrail'; import { SPACE_TYPE } from './spaces/config'; import { MAX_HYDRATION_URIS } from './search/constants'; import { createMeiliSink, meiliSinkBackendFromEnv } from './search/server/meili-sink'; +import { discoverableSql } from './search/server/discoverability'; // The `contrail` CLI (`pnpm backfill` / `contrail refresh`) fires `config.sinks` // on the backfill/refresh paths, so a fresh or re-synced install gets full search @@ -16,11 +17,11 @@ const searchSinks: NonNullable = ? [createMeiliSink(() => meiliSinkBackendFromEnv(process.env))] : []; -// Events hidden from discovery (preferences.showInDiscovery === false) are -// excluded; a missing field defaults to true so pre-existing records without -// `preferences` are included. Shared by every discovery-facing pipelineQuery. -const DISCOVERABLE_CONDITION = `(json_extract(r.record, '$.preferences.showInDiscovery') IS NULL - OR json_extract(r.record, '$.preferences.showInDiscovery') != 0)`; +// Events hidden from discovery (a falsey preferences.showInDiscovery) are +// excluded; a missing field defaults to discoverable so pre-existing records +// without `preferences` are included. The predicate is the single source of +// truth shared with the sink's in-memory filter and the geocode worklist. +const DISCOVERABLE_CONDITION = discoverableSql('r.record'); export const config: ContrailConfig = { namespace: 'rsvp.atmo', diff --git a/apps/web/src/lib/contrail.ts b/apps/web/src/lib/contrail.ts index 2900a88..b7f55af 100644 --- a/apps/web/src/lib/contrail.ts +++ b/apps/web/src/lib/contrail.ts @@ -280,9 +280,9 @@ export async function listEventRecordsFromContrail( /** * Hits the `listDiscoverable` pipelineQuery, which reuses the listRecords - * pipeline but adds a WHERE condition excluding events where - * `preferences.showInDiscovery === false`. Missing field is treated as true. - * Response shape is identical to listRecords. + * pipeline but adds a WHERE condition excluding events with a falsey + * `preferences.showInDiscovery` (false or 0). Missing field is treated as + * discoverable. Response shape is identical to listRecords. */ export async function listDiscoverableEventsFromContrail( client: Client, @@ -351,24 +351,17 @@ export async function listDiscoverableEventsByUrisFromContrail( */ export async function listConferenceTalksFromContrail( client: Client, - { - parentUri, - actor, - limit = 300 - }: { parentUri: string; actor?: ActorIdentifier; limit?: number } + { parentUri, actor, limit = 300 }: { parentUri: string; actor?: ActorIdentifier; limit?: number } ): Promise { - const response = await client.get( - 'rsvp.atmo.event.listTalks' as 'rsvp.atmo.event.listRecords', - { - params: { - ...(actor ? { actor } : {}), - parentUri, - sort: 'startsAt', - order: 'asc', - limit - } as ListEventsParams - } - ); + const response = await client.get('rsvp.atmo.event.listTalks' as 'rsvp.atmo.event.listRecords', { + params: { + ...(actor ? { actor } : {}), + parentUri, + sort: 'startsAt', + order: 'asc', + limit + } as ListEventsParams + }); if (!response.ok) return null; return response.data; diff --git a/apps/web/src/lib/contrail/index.ts b/apps/web/src/lib/contrail/index.ts index 41b44d8..acc5abf 100644 --- a/apps/web/src/lib/contrail/index.ts +++ b/apps/web/src/lib/contrail/index.ts @@ -27,16 +27,28 @@ if (!spacesAvailable()) { // which is fine: reads never ingest, so the sink never fires there. let searchSinkBackend: MeiliSinkBackend | null = null; +// The geocode cache lives in D1; the sink reads it (read-only, best-effort) to +// reproduce the _geo the external geocode job writes, so a live update doesn't +// drop it. Like searchSinkBackend, it's a module-level holder the env-bearing +// ensureInit populates — the sink no-ops on it until then. +let geocodeCacheDb: D1Database | null = null; + export const contrail = new Contrail({ ...config, ...(spaces ? { spaces } : {}), - sinks: [createMeiliSink(() => searchSinkBackend)] + sinks: [ + createMeiliSink( + () => searchSinkBackend, + () => geocodeCacheDb + ) + ] }); let initialized = false; let sinkConfigured = false; export async function ensureInit(db: D1Database, env?: MeiliSinkEnv) { + geocodeCacheDb = db; if (!initialized) { await contrail.init(db); initialized = true; diff --git a/apps/web/src/lib/search/server/address-norm.test.ts b/apps/web/src/lib/search/server/address-norm.test.ts new file mode 100644 index 0000000..fdb512d --- /dev/null +++ b/apps/web/src/lib/search/server/address-norm.test.ts @@ -0,0 +1,115 @@ +// apps/web/src/lib/search/server/address-norm.test.ts +import { describe, it, expect } from 'vitest'; +import { normalizeAddress, addressLocation, ADDRESS_TYPE } from './address-norm'; +import { addressNeedingGeocode } from './geocode-cache'; + +describe('normalizeAddress', () => { + it('joins all present fields in fixed order with | separators', () => { + const key = normalizeAddress({ + country: 'US', + locality: 'Dayton', + street: '905 East 3rd Street', + region: 'Ohio', + postalCode: '45402', + name: 'The Venue' + }); + // fixed order: name, street, locality, region, postalCode, country + expect(key).toBe('the venue|905 east 3rd street|dayton|ohio|45402|us'); + }); + + it('omits absent/empty fields entirely', () => { + expect(normalizeAddress({ locality: 'Dayton', country: 'US' })).toBe('dayton|us'); + }); + + it('lowercases, NFC-normalizes, collapses whitespace, trims edge punctuation', () => { + expect(normalizeAddress({ locality: 'Dayton,', region: ' Ohio State ' })).toBe( + 'dayton|ohio state' + ); + }); + + it('strips the separator char from field content', () => { + expect(normalizeAddress({ name: 'A|B', country: 'US' })).toBe('a b|us'); + }); + + it('returns null when nothing is present', () => { + expect(normalizeAddress({})).toBeNull(); + expect(normalizeAddress({ country: ' ' })).toBeNull(); + }); + + it('produces an identical key regardless of field insertion order', () => { + const a = normalizeAddress({ country: 'US', locality: 'Dayton' }); + const b = normalizeAddress({ locality: 'Dayton', country: 'US' }); + expect(a).toBe(b); + }); +}); + +describe('addressLocation', () => { + it('returns the first .address location', () => { + const loc = addressLocation({ + locations: [ + { $type: 'community.lexicon.location.geo', latitude: '1', longitude: '2' }, + { $type: ADDRESS_TYPE, locality: 'Dayton', country: 'US' } + ] + }); + expect(loc).toMatchObject({ locality: 'Dayton', country: 'US' }); + }); + + it('returns null when no address location is present', () => { + expect( + addressLocation({ locations: [{ $type: 'community.lexicon.location.geo' }] }) + ).toBeNull(); + expect(addressLocation({})).toBeNull(); + }); +}); + +// INV-A: the geocode job (write path) keys geocode_cache by +// normalizeAddress(addressNeedingGeocode(record)); the sink (read path) looks it +// back up by normalizeAddress(addressLocation(record)). If those two keys ever +// diverge, the _geo the job writes is invisible to the sink and near-me silently +// loses coordinates. These lock the round-trip - including across the Unicode/ +// punctuation/whitespace representation drift a PDS re-serialization can introduce. +describe('INV-A job <-> sink cache-key parity', () => { + const jobKey = (record: Record) => { + const loc = addressNeedingGeocode(record); + return loc ? normalizeAddress(loc) : null; + }; + const sinkKey = (record: Record) => { + const loc = addressLocation(record); + return loc ? normalizeAddress(loc) : null; + }; + + it('derives the same non-null key on both paths for one address-only record', () => { + const record = { + locations: [{ $type: ADDRESS_TYPE, name: 'The Venue', locality: 'Dayton', country: 'US' }] + }; + const k = jobKey(record); + expect(k).not.toBeNull(); + expect(sinkKey(record)).toBe(k); + }); + + it('matches across NFD/NFC, case, whitespace and edge-punctuation drift', () => { + // Precomposed (NFC) base strings as the source file stores them. + const venueNFC = 'Café Bar'; + const cityNFC = 'Zürich'; + // Decompose to NFD at runtime so the job side genuinely carries combining + // marks - independent of how this file happens to be encoded - then add the + // messy edges (doubled spaces, trailing comma) the job might first see. + const jobName = ` ${venueNFC.normalize('NFD')} `; + const jobCity = `${cityNFC.normalize('NFD')},`; + const jobRecord = { + locations: [{ $type: ADDRESS_TYPE, name: jobName, locality: jobCity, country: 'CH' }] + }; + // The same venue re-serialized: precomposed (NFC), already trimmed/cased. + const sinkRecord = { + locations: [{ $type: ADDRESS_TYPE, name: venueNFC, locality: cityNFC, country: 'ch' }] + }; + // The raw job/sink inputs differ in codepoints (NFD carries combining marks + // NFC lacks); only NFC folding makes them equal, so this fails if the + // normalize step ever drops it. + expect(jobName).not.toBe(venueNFC); + const expected = `${venueNFC.normalize('NFC').toLowerCase()}|${cityNFC.normalize('NFC').toLowerCase()}|ch`; + const k = jobKey(jobRecord); + expect(k).toBe(expected); + expect(sinkKey(sinkRecord)).toBe(k); + }); +}); diff --git a/apps/web/src/lib/search/server/address-norm.ts b/apps/web/src/lib/search/server/address-norm.ts new file mode 100644 index 0000000..1933eab --- /dev/null +++ b/apps/web/src/lib/search/server/address-norm.ts @@ -0,0 +1,41 @@ +// Deterministic cache key for a community.lexicon.location.address, shared +// byte-for-byte by the search sink (read path) and the external geocode job +// (write path). If these two ever diverge the cache keys stop matching and the +// _geo the job writes is invisible to the sink, so this module is the single +// source of truth for both. + +export const ADDRESS_TYPE = 'community.lexicon.location.address'; + +const SEP = '|'; +// Fixed field order, independent of the order fields appear in the record, so +// identical data always yields the same key. `name` MUST be included: ~750 +// events carry only a venue name + country, and dropping it would collapse +// every venue-only event in a country into one row and cross-contaminate coords. +const FIELDS = ['name', 'street', 'locality', 'region', 'postalCode', 'country']; + +export function normalizeAddress(loc: Record): string | null { + const key = FIELDS.map((f) => (typeof loc[f] === 'string' ? (loc[f] as string) : '')) + .map((v) => + v + .normalize('NFC') + .toLowerCase() + .replace(/\s+/g, ' ') + .replace(/\|/g, ' ') + .replace(/^[\s,.;:/-]+|[\s,.;:/-]+$/g, '') + .trim() + ) + .filter((v) => v !== '') + .join(SEP); + return key === '' ? null : key; +} + +/** First community.lexicon.location.address in a record's locations[], or null. */ +export function addressLocation(record: Record): Record | null { + const locs = Array.isArray(record?.locations) ? record.locations : []; + for (const l of locs) { + if (l && typeof l === 'object' && (l as Record).$type === ADDRESS_TYPE) { + return l as Record; + } + } + return null; +} diff --git a/apps/web/src/lib/search/server/d1-http.test.ts b/apps/web/src/lib/search/server/d1-http.test.ts new file mode 100644 index 0000000..3c6e3df --- /dev/null +++ b/apps/web/src/lib/search/server/d1-http.test.ts @@ -0,0 +1,61 @@ +import { describe, it, expect, vi } from 'vitest'; +import { createD1Client } from './d1-http'; + +const CFG = { accountId: 'acct', databaseId: 'db', apiToken: 'tok' }; + +function fakeFetch(body: unknown, status = 200) { + const calls: { url: string; method: string; auth: string; body: unknown }[] = []; + const fn = vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + calls.push({ + url: String(input), + method: init?.method ?? 'GET', + auth: (init?.headers as Record)?.authorization ?? '', + body: init?.body ? JSON.parse(String(init.body)) : undefined + }); + return new Response(JSON.stringify(body), { status }); + }); + return { fn: fn as unknown as typeof fetch, calls }; +} + +describe('createD1Client.query', () => { + it('POSTs sql+params to the D1 query endpoint and returns the first result set', async () => { + const { fn, calls } = fakeFetch({ + success: true, + result: [{ results: [{ uri: 'at://x' }], success: true }] + }); + const rows = await createD1Client(CFG, fn).query( + 'SELECT uri FROM records_event WHERE did = ?', + ['did:plc:a'] + ); + expect(rows).toEqual([{ uri: 'at://x' }]); + expect(calls[0].url).toBe( + 'https://api.cloudflare.com/client/v4/accounts/acct/d1/database/db/query' + ); + expect(calls[0].method).toBe('POST'); + expect(calls[0].auth).toBe('Bearer tok'); + expect(calls[0].body).toEqual({ + sql: 'SELECT uri FROM records_event WHERE did = ?', + params: ['did:plc:a'] + }); + }); + + it('returns [] when the result set is empty', async () => { + const { fn } = fakeFetch({ success: true, result: [{ results: [], success: true }] }); + expect(await createD1Client(CFG, fn).query('SELECT 1')).toEqual([]); + }); + + it('throws on an HTTP error', async () => { + const { fn } = fakeFetch({}, 401); + await expect(createD1Client(CFG, fn).query('SELECT 1')).rejects.toThrow(/401/); + }); + + it('throws on a D1 error envelope (success:false)', async () => { + const { fn } = fakeFetch({ success: false, errors: [{ message: 'bad sql' }] }); + await expect(createD1Client(CFG, fn).query('SELECT')).rejects.toThrow(/bad sql/); + }); + + it('falls back to "unknown D1 error" when all error objects lack a message', async () => { + const { fn } = fakeFetch({ success: false, errors: [{}, {}] }); + await expect(createD1Client(CFG, fn).query('SELECT 1')).rejects.toThrow(/unknown D1 error/); + }); +}); diff --git a/apps/web/src/lib/search/server/d1-http.ts b/apps/web/src/lib/search/server/d1-http.ts new file mode 100644 index 0000000..195f914 --- /dev/null +++ b/apps/web/src/lib/search/server/d1-http.ts @@ -0,0 +1,43 @@ +// Thin Cloudflare D1 REST client for the external geocode job, which runs off +// Cloudflare (no Worker D1 binding). One parameterized-query method; the job +// does all its reads/writes through it. +export interface D1HttpConfig { + accountId: string; + databaseId: string; + apiToken: string; +} + +export interface D1Client { + query>(sql: string, params?: unknown[]): Promise; +} + +export function createD1Client(cfg: D1HttpConfig, fetchImpl: typeof fetch = fetch): D1Client { + const endpoint = `https://api.cloudflare.com/client/v4/accounts/${cfg.accountId}/d1/database/${cfg.databaseId}/query`; + return { + async query>(sql: string, params: unknown[] = []): Promise { + const res = await fetchImpl(endpoint, { + method: 'POST', + headers: { + authorization: `Bearer ${cfg.apiToken}`, + 'content-type': 'application/json' + }, + body: JSON.stringify({ sql, params }) + }); + if (!res.ok) throw new Error(`D1 query failed: ${res.status}`); + const body = (await res.json()) as { + success: boolean; + result?: Array<{ results: T[] }>; + errors?: Array<{ message?: string }>; + }; + if (!body.success) { + const msg = + body.errors + ?.map((e) => e.message) + .filter(Boolean) + .join('; ') || 'unknown D1 error'; + throw new Error(`D1 query error: ${msg}`); + } + return body.result?.[0]?.results ?? []; + } + }; +} diff --git a/apps/web/src/lib/search/server/discoverability.test.ts b/apps/web/src/lib/search/server/discoverability.test.ts new file mode 100644 index 0000000..5530dca --- /dev/null +++ b/apps/web/src/lib/search/server/discoverability.test.ts @@ -0,0 +1,60 @@ +import { describe, it, expect } from 'vitest'; +import { discoverableSql, isHiddenFromDiscovery } from './discoverability'; + +// The single rule both the SQL read paths (contrail.config.ts pipelineQueries, +// the geocode worklist) and the in-memory sink filter (meili-sink.ts) must agree +// on. If they drift, the index and the hydrated read path diverge: the sink +// indexes an event D1 hides (or vice-versa), creating a phantom doc. +// +// `sqlVisible` is the ground-truth verdict of the shared SQL predicate +// json_extract(...,'$.preferences.showInDiscovery') IS NULL OR != 0 +// evaluated against REAL SQLite (node:sqlite). Each row was confirmed by running +// the exact predicate over the JSON value; the table encodes those results so the +// parity check stays inside the Worker-typed test env (no node:sqlite/@types/node +// import here). SQLite maps JSON false→0 and true→1 and passes numbers through, so +// `!= 0` hides BOTH boolean false AND numeric 0 (and 0.0); strings, null, true, +// and missing all stay visible. +const PARITY: { label: string; record: Record; sqlVisible: boolean }[] = [ + { label: 'true', record: { preferences: { showInDiscovery: true } }, sqlVisible: true }, + { label: 'false', record: { preferences: { showInDiscovery: false } }, sqlVisible: false }, + { label: 'numeric 0', record: { preferences: { showInDiscovery: 0 } }, sqlVisible: false }, + { label: 'float 0.0', record: { preferences: { showInDiscovery: 0.0 } }, sqlVisible: false }, + { label: 'numeric 1', record: { preferences: { showInDiscovery: 1 } }, sqlVisible: true }, + { label: 'explicit null', record: { preferences: { showInDiscovery: null } }, sqlVisible: true }, + { + label: 'string "false"', + record: { preferences: { showInDiscovery: 'false' } }, + sqlVisible: true + }, + { label: 'string "0"', record: { preferences: { showInDiscovery: '0' } }, sqlVisible: true }, + { label: 'field absent', record: { preferences: {} }, sqlVisible: true }, + { label: 'preferences absent', record: {}, sqlVisible: true } +]; + +describe('isHiddenFromDiscovery', () => { + for (const { label, record, sqlVisible } of PARITY) { + it(`mirrors the SQL verdict for showInDiscovery=${label}`, () => { + expect(isHiddenFromDiscovery(record)).toBe(!sqlVisible); + }); + } + + it('tolerates a null/non-object record without throwing', () => { + expect(isHiddenFromDiscovery(null as unknown as Record)).toBe(false); + expect(isHiddenFromDiscovery({ preferences: null } as unknown as Record)).toBe( + false + ); + }); +}); + +describe('discoverableSql', () => { + it('emits the != 0 predicate over the given record column', () => { + const sql = discoverableSql('r.record'); + expect(sql).toContain("json_extract(r.record, '$.preferences.showInDiscovery')"); + expect(sql).toContain('IS NULL'); + expect(sql).toContain('!= 0'); + }); + + it('parameterizes the column so different table aliases reuse one rule', () => { + expect(discoverableSql('x.record')).toContain('json_extract(x.record'); + }); +}); diff --git a/apps/web/src/lib/search/server/discoverability.ts b/apps/web/src/lib/search/server/discoverability.ts new file mode 100644 index 0000000..6a9cf66 --- /dev/null +++ b/apps/web/src/lib/search/server/discoverability.ts @@ -0,0 +1,35 @@ +// Single source of truth for "is this event discoverable?". Used by BOTH the SQL +// read paths (contrail.config.ts pipelineQueries and the external geocode +// worklist) AND the in-memory sink filter (meili-sink.ts). The two MUST agree on +// semantics or the search index and the hydrated read path drift: an event the +// sink indexes but D1 hides (or the reverse) becomes a phantom, present in Meili +// yet dropped at hydration, or missing from search yet listed by D1. +// +// The rule: preferences.showInDiscovery hides the event when it is FALSEY in the +// JSON sense (boolean false OR numeric 0, since SQLite stores JSON false as 0). A +// missing/null/true value, any nonzero number, or any string stays discoverable, +// so pre-existing records without `preferences` are included by default. The SQL +// form uses `!= 0`; the JS form mirrors it with `=== false || === 0` (JS treats +// 0.0 as 0, matching SQLite). Parity across the full value surface is locked by +// discoverability.test.ts against real-SQLite ground truth. + +const PREF_PATH = '$.preferences.showInDiscovery'; + +/** SQL predicate (true = discoverable) over a record column expression, e.g. + * `discoverableSql('r.record')`. The one definition both the contrail + * pipelineQueries and the geocode worklist import, so the two SQL sites cannot + * drift from each other or from the JS mirror below. */ +export function discoverableSql(recordCol: string): string { + return `(json_extract(${recordCol}, '${PREF_PATH}') IS NULL + OR json_extract(${recordCol}, '${PREF_PATH}') != 0)`; +} + +/** In-memory mirror of discoverableSql for the sink: true when the author hid the + * event from discovery. Hides on a falsey showInDiscovery (boolean false OR + * numeric 0, including 0.0) to match SQLite's `!= 0`; missing/null/true/1/strings + * stay discoverable. Tolerates a null or non-object record/preferences. */ +export function isHiddenFromDiscovery(record: Record): boolean { + const prefs = record?.preferences as { showInDiscovery?: unknown } | null | undefined; + const v = prefs?.showInDiscovery; + return v === false || v === 0; +} diff --git a/apps/web/src/lib/search/server/geocode-cache.test.ts b/apps/web/src/lib/search/server/geocode-cache.test.ts new file mode 100644 index 0000000..feaf271 --- /dev/null +++ b/apps/web/src/lib/search/server/geocode-cache.test.ts @@ -0,0 +1,128 @@ +import { describe, it, expect } from 'vitest'; +import { + backoffMs, + isEligible, + groupEventsByNorm, + addressNeedingGeocode, + MAX_FAIL +} from './geocode-cache'; +import type { GeocodeCacheRow } from './geocode-cache'; +import { ADDRESS_TYPE } from './address-norm'; + +const DAY = 86_400_000; +const NOW = 1_000 * DAY; + +function negative(failCount: number, ageMs: number): GeocodeCacheRow { + return { + address_norm: 'x', + lat: null, + lng: null, + precision: null, + source: 'locationiq', + geocoded_at: NOW - ageMs, + fail_count: failCount, + last_error: 'not found' + }; +} + +describe('backoffMs', () => { + it('escalates 1d / 7d / 30d', () => { + expect(backoffMs(1)).toBe(DAY); + expect(backoffMs(2)).toBe(7 * DAY); + expect(backoffMs(3)).toBe(30 * DAY); + }); +}); + +describe('isEligible', () => { + it('treats an absent row as work', () => { + expect(isEligible(undefined, NOW)).toBe(true); + }); + + it('never re-geocodes a resolved row', () => { + const resolved: GeocodeCacheRow = { + address_norm: 'x', + lat: 50, + lng: 4, + precision: 'locality', + source: 'locationiq', + geocoded_at: NOW - 999 * DAY, + fail_count: 0, + last_error: null + }; + expect(isEligible(resolved, NOW)).toBe(false); + }); + + it('retries a negative row only after its backoff elapses', () => { + expect(isEligible(negative(1, 0.5 * DAY), NOW)).toBe(false); // < 1d + expect(isEligible(negative(1, 1.5 * DAY), NOW)).toBe(true); // >= 1d + expect(isEligible(negative(2, 3 * DAY), NOW)).toBe(false); // < 7d + expect(isEligible(negative(2, 8 * DAY), NOW)).toBe(true); // >= 7d + }); + + it('hard-stops at fail_count >= MAX_FAIL', () => { + expect(isEligible(negative(MAX_FAIL, 999 * DAY), NOW)).toBe(false); + }); + + it('--retry-negative ignores backoff and the hard stop', () => { + expect(isEligible(negative(MAX_FAIL, 0), NOW, true)).toBe(true); + }); +}); + +describe('addressNeedingGeocode', () => { + const addr = { $type: ADDRESS_TYPE, locality: 'Dayton', country: 'US' }; + const geo = (lat: string, lng: string) => ({ + $type: 'community.lexicon.location.geo', + latitude: lat, + longitude: lng + }); + const fsq = (lat: string, lng: string) => ({ + $type: 'community.lexicon.location.fsq', + latitude: lat, + longitude: lng + }); + + it('returns the address location for an address-only event', () => { + expect(addressNeedingGeocode({ locations: [addr] })).toMatchObject({ + locality: 'Dayton', + country: 'US' + }); + }); + + it('returns null when a geo location already resolves to coordinates (no overwrite)', () => { + expect(addressNeedingGeocode({ locations: [geo('40', '-105'), addr] })).toBeNull(); + }); + + it('returns null when an fsq location already resolves to coordinates (F1a)', () => { + expect(addressNeedingGeocode({ locations: [fsq('50.8', '4.3'), addr] })).toBeNull(); + }); + + it('returns the address location when the only coordinate location is out of range (F1b)', () => { + expect(addressNeedingGeocode({ locations: [geo('999', '0'), addr] })).toMatchObject({ + locality: 'Dayton' + }); + }); + + it('returns null when there is no address location to geocode', () => { + expect(addressNeedingGeocode({ locations: [geo('40', '-105')] })).toBeNull(); + expect(addressNeedingGeocode({ locations: [] })).toBeNull(); + }); +}); + +describe('groupEventsByNorm', () => { + it('buckets events by normalized address and skips unkeyable ones', () => { + const ev = (rkey: string, loc: Record) => ({ + uri: `at://did:plc:a/community.lexicon.calendar.event/${rkey}`, + did: 'did:plc:a', + rkey, + loc: { $type: ADDRESS_TYPE, ...loc } + }); + const map = groupEventsByNorm([ + ev('1', { locality: 'Dayton', country: 'US' }), + ev('2', { locality: 'Dayton', country: 'US' }), + ev('3', { locality: 'Berlin', country: 'DE' }), + ev('4', {}) // unkeyable -> skipped + ]); + expect([...map.keys()].sort()).toEqual(['berlin|de', 'dayton|us']); + expect(map.get('dayton|us')!.map((e: { rkey: string }) => e.rkey)).toEqual(['1', '2']); + }); +}); diff --git a/apps/web/src/lib/search/server/geocode-cache.ts b/apps/web/src/lib/search/server/geocode-cache.ts new file mode 100644 index 0000000..c883fe1 --- /dev/null +++ b/apps/web/src/lib/search/server/geocode-cache.ts @@ -0,0 +1,80 @@ +// Pure decision logic for the geocode_cache worklist: the row shape, the +// negative-cache retry/backoff policy, and grouping worklist events by their +// normalized address. No I/O — the job feeds rows in and acts on the verdicts. +import { normalizeAddress, addressLocation } from './address-norm'; +import { recordGeo } from './normalize'; + +export interface GeocodeCacheRow { + address_norm: string; + lat: number | null; + lng: number | null; + precision: string | null; + source: string | null; + geocoded_at: number; // epoch ms + fail_count: number; + last_error: string | null; +} + +export interface WorklistEvent { + uri: string; + did: string; + rkey: string; + loc: Record; +} + +const DAY = 86_400_000; +/** No more automatic attempts once fail_count reaches this (~4 tries / ~5 wks). */ +export const MAX_FAIL = 4; + +/** No-match backoff: eligible again after 1d (fail 1), 7d (2), 30d (3+). */ +export function backoffMs(failCount: number): number { + if (failCount <= 1) return DAY; + if (failCount === 2) return 7 * DAY; + return 30 * DAY; +} + +/** Is this address work this run? Absent → yes. Resolved → no. Negative → + * retryable iff fail_count < MAX_FAIL and the backoff elapsed (or forced). */ +export function isEligible( + row: GeocodeCacheRow | undefined, + now: number, + retryNegative = false +): boolean { + if (!row) return true; + if (row.lat !== null && row.lng !== null) return false; // resolved, done + if (retryNegative) return true; + if (row.fail_count >= MAX_FAIL) return false; + return now - row.geocoded_at >= backoffMs(row.fail_count); +} + +/** The address location to geocode for a record, or null when there's nothing + * to do. Returns null if the record ALREADY resolves to coordinates the index + * derives (recordGeo: geo/fsq/hthree, in-range) — geocoding it would overwrite + * precise coords with approximate ones — and null if it carries no address at + * all. So an event with an address plus an out-of-range geo/hthree location + * (recordGeo undefined) correctly still gets its address geocoded. This is the + * in-memory worklist filter that keeps "needs geocoding" aligned with the + * sink's _geo derivation, since the SQL worklist can only coarsely pre-filter. */ +export function addressNeedingGeocode( + record: Record +): Record | null { + if (recordGeo(record)) return null; + return addressLocation(record); +} + +/** Bucket worklist events by their normalized address; events that don't + * normalize to a key (e.g. empty address) are dropped. */ +export function groupEventsByNorm(events: WorklistEvent[]): Map { + const map = new Map(); + for (const e of events) { + const norm = normalizeAddress(e.loc); + if (!norm) continue; + let arr = map.get(norm); + if (!arr) { + arr = []; + map.set(norm, arr); + } + arr.push(e); + } + return map; +} diff --git a/apps/web/src/lib/search/server/geocoder.test.ts b/apps/web/src/lib/search/server/geocoder.test.ts new file mode 100644 index 0000000..c11ba58 --- /dev/null +++ b/apps/web/src/lib/search/server/geocoder.test.ts @@ -0,0 +1,140 @@ +// apps/web/src/lib/search/server/geocoder.test.ts +import { describe, it, expect, vi } from 'vitest'; +import { + addressToQuery, + derivePrecision, + createGeocoder, + requireGeocoderForBulk +} from './geocoder'; + +describe('addressToQuery', () => { + it('joins present fields in fixed order with commas, preserving original case', () => { + expect( + addressToQuery({ + country: 'België / Belgique / Belgien', + street: 'Cantersteen 41', + locality: 'Bruxelles - Brussel', + region: 'Brussel-Hoofdstad' + }) + ).toBe('Cantersteen 41, Bruxelles - Brussel, Brussel-Hoofdstad, België / Belgique / Belgien'); + }); + + it('skips absent/blank fields', () => { + expect(addressToQuery({ locality: 'Dayton', region: ' ', country: 'US' })).toBe('Dayton, US'); + }); +}); + +describe('derivePrecision', () => { + it('classifies a house/building hit as rooftop', () => { + expect(derivePrecision({ addresstype: 'house', place_rank: 30 })).toBe('rooftop'); + }); + it('classifies a road hit as street', () => { + expect(derivePrecision({ type: 'road', class: 'highway' })).toBe('street'); + }); + it('classifies a city/region hit as locality', () => { + expect(derivePrecision({ addresstype: 'city' })).toBe('locality'); + expect(derivePrecision({ type: 'administrative', class: 'boundary' })).toBe('locality'); + }); + it('falls back to the raw type when unclassifiable', () => { + expect(derivePrecision({ type: 'attraction' })).toBe('attraction'); + expect(derivePrecision({})).toBe('unknown'); + }); +}); + +describe('createGeocoder', () => { + function fakeFetch(body: unknown, status = 200) { + const calls: { url: string; headers: Record }[] = []; + const fn = vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + calls.push({ url: String(input), headers: (init?.headers ?? {}) as Record }); + return new Response(JSON.stringify(body), { status }); + }); + return { fn: fn as unknown as typeof fetch, calls }; + } + + it('defaults to public Nominatim with the atmo user-agent', async () => { + const { fn, calls } = fakeFetch([{ lat: '50.84', lon: '4.36', addresstype: 'city' }]); + const geo = createGeocoder({}, fn); + const point = await geo.geocode('Bruxelles, BE'); + expect(point).toEqual({ lat: 50.84, lng: 4.36, precision: 'locality' }); + expect(calls[0].url).toContain('https://nominatim.openstreetmap.org/search'); + expect(calls[0].url).toContain('q=Bruxelles%2C+BE'); + expect(calls[0].headers['user-agent']).toContain('atmo-events'); + }); + + it('switches to a keyed LocationIQ endpoint when GEOCODER_KEY+URL are set', async () => { + const { fn, calls } = fakeFetch([{ lat: '52.5', lon: '13.4', addresstype: 'road' }]); + const geo = createGeocoder( + { GEOCODER_URL: 'https://us1.locationiq.com/v1/search', GEOCODER_KEY: 'tok' }, + fn + ); + await geo.geocode('Berlin'); + expect(calls[0].url).toContain('https://us1.locationiq.com/v1/search'); + expect(calls[0].url).toContain('key=tok'); + }); + + it('returns null on an empty result set (no-match)', async () => { + const { fn } = fakeFetch([]); + expect(await createGeocoder({}, fn).geocode('nowhere')).toBeNull(); + }); + + it('returns null on a 404 (LocationIQ no-match) so it negative-caches, not retries', async () => { + const { fn } = fakeFetch({ error: 'Unable to geocode' }, 404); + expect(await createGeocoder({}, fn).geocode('nowhere')).toBeNull(); + }); + + it('throws on a transient HTTP error (429/5xx) so the caller retries', async () => { + const { fn } = fakeFetch({}, 429); + await expect(createGeocoder({}, fn).geocode('x')).rejects.toThrow(/429/); + }); +}); + +describe('requireGeocoderForBulk', () => { + it('allows any run when a geocoder key is set', () => { + expect(() => + requireGeocoderForBulk({ hasKey: true, dryRun: false, limit: 0, allowPublic: false }) + ).not.toThrow(); + expect(() => + requireGeocoderForBulk({ hasKey: true, dryRun: false, limit: 5000, allowPublic: false }) + ).not.toThrow(); + }); + + it('allows a keyless small drip (1..25)', () => { + expect(() => + requireGeocoderForBulk({ hasKey: false, dryRun: false, limit: 25, allowPublic: false }) + ).not.toThrow(); + expect(() => + requireGeocoderForBulk({ hasKey: false, dryRun: false, limit: 1, allowPublic: false }) + ).not.toThrow(); + }); + + it('blocks a keyless run over the drip ceiling', () => { + expect(() => + requireGeocoderForBulk({ hasKey: false, dryRun: false, limit: 26, allowPublic: false }) + ).toThrow(/LocationIQ|allow-public-nominatim/); + }); + + it('blocks a keyless uncapped run (limit 0 = no cap, the worst case)', () => { + expect(() => + requireGeocoderForBulk({ hasKey: false, dryRun: false, limit: 0, allowPublic: false }) + ).toThrow(/LocationIQ|allow-public-nominatim/); + }); + + it('blocks a keyless negative limit (uncapped, not a tiny drip)', () => { + // A stray `--limit -1` must not masquerade as a 1-call drip and slip the gate. + expect(() => + requireGeocoderForBulk({ hasKey: false, dryRun: false, limit: -1, allowPublic: false }) + ).toThrow(/LocationIQ|allow-public-nominatim/); + }); + + it('lets --allow-public-nominatim override a bulk keyless run', () => { + expect(() => + requireGeocoderForBulk({ hasKey: false, dryRun: false, limit: 0, allowPublic: true }) + ).not.toThrow(); + }); + + it('never blocks a dry run (it makes no geocoder calls)', () => { + expect(() => + requireGeocoderForBulk({ hasKey: false, dryRun: true, limit: 0, allowPublic: false }) + ).not.toThrow(); + }); +}); diff --git a/apps/web/src/lib/search/server/geocoder.ts b/apps/web/src/lib/search/server/geocoder.ts new file mode 100644 index 0000000..c752e2b --- /dev/null +++ b/apps/web/src/lib/search/server/geocoder.ts @@ -0,0 +1,144 @@ +// apps/web/src/lib/search/server/geocoder.ts +// One Nominatim-compatible geocoding client, config-selected. LocationIQ is +// API-compatible with Nominatim (same /search?q=&format=&limit= request, same +// lat/lon/type/class/display_name response — it IS hosted Nominatim), so this +// is one implementation parameterized by endpoint + optional key, not two +// behind an interface. Default = public Nominatim (parity with atmo today); +// set GEOCODER_KEY (+ GEOCODER_URL) to use LocationIQ. Used only by the +// external geocode job — geocoding never runs on the Worker hot path. + +export interface GeoPoint { + lat: number; + lng: number; + /** Provider-reported granularity tier (rooftop/street/locality/...). */ + precision?: string; +} + +export interface Geocoder { + geocode(q: string): Promise; +} + +export interface GeocoderEnv { + GEOCODER_URL?: string; + GEOCODER_KEY?: string; + GEOCODER_USER_AGENT?: string; +} + +const DEFAULT_URL = 'https://nominatim.openstreetmap.org/search'; +const DEFAULT_USER_AGENT = 'atmo-events (https://atmo.rsvp)'; + +// Same fixed order as the cache key, but a human-readable freeform query (commas, +// original case/diacritics) — both Nominatim and LocationIQ handle UTF-8 / non- +// Latin scripts. Freeform tolerates the messy data (city-in-name, JP street-only, +// trailing punctuation) a structured per-field query would drop. +const QUERY_FIELDS = ['name', 'street', 'locality', 'region', 'postalCode', 'country']; + +/** Largest keyless run we treat as a sanctioned "small drip" against public + * Nominatim. Above this (or uncapped), a backfill must use LocationIQ. */ +export const PUBLIC_NOMINATIM_DRIP_MAX = 25; + +/** Enforce the file's contract that a BULK backfill uses LocationIQ, not public + * Nominatim (whose usage policy a large unkeyed run would breach, risking a + * silent IP ban). Throws for a keyless, non-dry-run, non-overridden run that is + * uncapped (limit <= 0) or over the drip ceiling; a small keyless drip stays + * allowed. The hazard is request VOLUME, so the gate is on size, not on the + * mere use of Nominatim. limit <= 0 (not just === 0) counts as uncapped so a + * stray negative can't masquerade as a tiny drip and slip the gate. */ +export function requireGeocoderForBulk(opts: { + hasKey: boolean; + dryRun: boolean; + limit: number; + allowPublic: boolean; +}): void { + if (opts.hasKey || opts.dryRun || opts.allowPublic) return; + const bulk = opts.limit <= 0 || opts.limit > PUBLIC_NOMINATIM_DRIP_MAX; + if (bulk) { + throw new Error( + `Keyless public Nominatim is only allowed for a small drip (--limit 1..${PUBLIC_NOMINATIM_DRIP_MAX}). ` + + 'Set GEOCODER_KEY (LocationIQ) for a bulk/uncapped backfill, or pass --allow-public-nominatim to override.' + ); + } +} + +export function addressToQuery(loc: Record): string { + return QUERY_FIELDS.map((f) => (typeof loc[f] === 'string' ? (loc[f] as string).trim() : '')) + .filter((v) => v !== '') + .join(', '); +} + +type NominatimHit = { + lat?: string; + lon?: string; + display_name?: string; + type?: string; + class?: string; + place_rank?: number; + addresstype?: string; +}; + +const ROOFTOP = new Set(['house', 'building', 'address', 'house_number']); +const STREET = new Set(['road', 'street', 'residential', 'pedestrian']); +const LOCALITY = new Set([ + 'city', + 'town', + 'village', + 'hamlet', + 'suburb', + 'locality', + 'municipality', + 'administrative', + 'state', + 'region', + 'province', + 'country' +]); + +export function derivePrecision(top: { + type?: string; + class?: string; + place_rank?: number; + addresstype?: string; +}): string { + const t = top.addresstype || top.type || ''; + if (ROOFTOP.has(t) || (top.place_rank ?? 0) >= 30) return 'rooftop'; + if (STREET.has(t) || top.class === 'highway') return 'street'; + if (LOCALITY.has(t) || top.class === 'boundary' || top.class === 'place') return 'locality'; + return t || 'unknown'; +} + +export function createGeocoder(env: GeocoderEnv = {}, fetchImpl: typeof fetch = fetch): Geocoder { + const base = env.GEOCODER_URL || DEFAULT_URL; + const key = env.GEOCODER_KEY; + const userAgent = env.GEOCODER_USER_AGENT || DEFAULT_USER_AGENT; + + return { + async geocode(q: string): Promise { + const url = new URL(base); + url.searchParams.set('q', q); + url.searchParams.set('format', 'json'); + url.searchParams.set('limit', '1'); + if (key) url.searchParams.set('key', key); + + const res = await fetchImpl(url, { + headers: { accept: 'application/json', 'user-agent': userAgent } + }); + // NO-MATCH vs TRANSIENT. Nominatim signals no-match as 200 + []; LocationIQ + // signals it as 404 (e.g. {"error":"Unable to geocode"}). Treat 404 as a + // no-match (return null → negative cache w/ backoff) so an ungeocodable + // address isn't retried every run. Other non-2xx (429 rate-limit, 5xx, + // transport) throw → the job treats them as TRANSIENT and retries next run. + if (res.status === 404) return null; + if (!res.ok) throw new Error(`geocode request failed: ${res.status}`); + + const results = (await res.json()) as NominatimHit[]; + const top = Array.isArray(results) ? results[0] : undefined; + if (!top) return null; + + const lat = Number(top.lat); + const lng = Number(top.lon); + if (!Number.isFinite(lat) || !Number.isFinite(lng)) return null; + + return { lat, lng, precision: derivePrecision(top) }; + } + }; +} diff --git a/apps/web/src/lib/search/server/meili-sink.integration.test.ts b/apps/web/src/lib/search/server/meili-sink.integration.test.ts index 791ce20..97723d6 100644 --- a/apps/web/src/lib/search/server/meili-sink.integration.test.ts +++ b/apps/web/src/lib/search/server/meili-sink.integration.test.ts @@ -8,8 +8,14 @@ // MEILI_TEST_URL=http://localhost:7700 MEILI_TEST_KEY=masterKey \ // pnpm vitest run src/lib/search/server/meili-sink.integration.test.ts import { describe, it, expect, beforeAll, afterAll } from 'vitest'; -import { createMeiliSink, applyMeiliSettings, type MeiliSinkBackend } from './meili-sink'; +import { + createMeiliSink, + applyMeiliSettings, + MeiliEventIndex, + type MeiliSinkBackend +} from './meili-sink'; import { searchEvents, nearMeEvents, type SearchBackend } from './meili'; +import { eventToSearchDoc } from './normalize'; const URL = process.env.MEILI_TEST_URL; const KEY = process.env.MEILI_TEST_KEY ?? 'masterKey'; @@ -101,6 +107,55 @@ run('MeiliSink ↔ read client, live against real Meilisearch', () => { expect(hit?.distanceMeters).toBeGreaterThanOrEqual(0); }); + it('the geocode job upsert attaches _geo to an address-only event, keeping its other fields', async () => { + const addrUri = 'at://did:plc:alice/community.lexicon.calendar.event/addr-only'; + const addrRecord = { + name: 'Antwerp Atproto Drinks', + description: 'address-only event', + startsAt: FUTURE, + locations: [ + { $type: 'community.lexicon.location.address', locality: 'Antwerp', country: 'BE' } + ] + }; + // An address-only event: indexed by the sink, but with no _geo yet. + await sink.onRecords([created(addrUri, addrRecord)], { phase: 'live' }); + await eventually( + () => searchEvents(readBackend, { q: 'Antwerp Atproto', limit: 10, offset: 0 }), + (r) => r.hits.some((h) => h.uri === addrUri) + ); + + // Attach coordinates the way the external geocode job does: rebuild the FULL + // doc (eventToSearchDoc) and upsert it with _geo — idempotent, stub-free, and + // identical to the doc the sink writes. + const doc = eventToSearchDoc({ + uri: addrUri, + did: 'did:plc:alice', + collection: EVENT, + rkey: addrUri.split('/').pop()!, + record: addrRecord + }); + doc._geo = { lat: 51.2194, lng: 4.4025 }; + await new MeiliEventIndex(backend).upsert([doc]); + + // near-me finds it (so _geo landed) AND text still matches (so name/startsAt + // survived) — both read-path filters pass against the full doc. + const near = await eventually( + () => + nearMeEvents(readBackend, { + lat: 51.2194, + lng: 4.4025, + radiusMeters: 5000, + limit: 10, + offset: 0 + }), + (r) => r.hits.some((h) => h.uri === addrUri) + ); + expect(near.hits.map((h) => h.uri)).toContain(addrUri); + + const text = await searchEvents(readBackend, { q: 'Antwerp Atproto', limit: 10, offset: 0 }); + expect(text.hits.map((h) => h.uri)).toContain(addrUri); // name survived the merge + }); + it('removes a deleted event so the read path no longer finds it', async () => { await sink.onRecords( [{ kind: 'deleted', uri, did: 'did:plc:alice', collection: EVENT, rkey: 'round-trip' }], diff --git a/apps/web/src/lib/search/server/meili-sink.test.ts b/apps/web/src/lib/search/server/meili-sink.test.ts index aa19357..32d1d00 100644 --- a/apps/web/src/lib/search/server/meili-sink.test.ts +++ b/apps/web/src/lib/search/server/meili-sink.test.ts @@ -3,9 +3,11 @@ import { createMeiliSink, meiliSinkBackendFromEnv, applyMeiliSettings, + MeiliEventIndex, type MeiliSinkBackend } from './meili-sink'; import { searchDocId } from './normalize'; +import { normalizeAddress, ADDRESS_TYPE } from './address-norm'; const BACKEND: MeiliSinkBackend = { url: 'http://meili.local', apiKey: 'admin-key' }; @@ -66,7 +68,11 @@ describe('meiliSinkBackendFromEnv', () => { describe('createMeiliSink onRecords', () => { it('upserts a created event as a normalized doc (PUT documents)', async () => { const { fn, calls } = fakeFetch(); - const sink = createMeiliSink(() => BACKEND, fn); + const sink = createMeiliSink( + () => BACKEND, + () => null, + fn + ); await sink.onRecords( [ @@ -102,7 +108,11 @@ describe('createMeiliSink onRecords', () => { it('removes a deleted event by derived id (delete-batch)', async () => { const { fn, calls } = fakeFetch(); - const sink = createMeiliSink(() => BACKEND, fn); + const sink = createMeiliSink( + () => BACKEND, + () => null, + fn + ); const uri = 'at://did:plc:alice/community.lexicon.calendar.event/2'; await sink.onRecords( @@ -118,7 +128,11 @@ describe('createMeiliSink onRecords', () => { it('removes (does not index) a created event hidden from discovery', async () => { const { fn, calls } = fakeFetch(); - const sink = createMeiliSink(() => BACKEND, fn); + const sink = createMeiliSink( + () => BACKEND, + () => null, + fn + ); const uri = 'at://did:plc:alice/community.lexicon.calendar.event/hidden'; await sink.onRecords( @@ -133,9 +147,34 @@ describe('createMeiliSink onRecords', () => { expect(del!.body).toEqual([searchDocId(uri)]); }); + it('removes a created event whose showInDiscovery is a numeric 0 (matches the SQL != 0 filter)', async () => { + const { fn, calls } = fakeFetch(); + const sink = createMeiliSink( + () => BACKEND, + () => null, + fn + ); + const uri = 'at://did:plc:alice/community.lexicon.calendar.event/zero'; + + await sink.onRecords( + [created(uri, { name: 'SerializedFalse', preferences: { showInDiscovery: 0 } })], + { phase: 'live' } + ); + + // Some serializers emit booleans as 0/1; the D1 worklist/read filter uses + // `!= 0`, so the sink must hide 0 too or it would index an event D1 drops. + expect(calls.find((c) => c.method === 'PUT')).toBeUndefined(); + const del = calls.find((c) => c.url.endsWith('/documents/delete-batch')); + expect(del!.body).toEqual([searchDocId(uri)]); + }); + it('indexes a created event when showInDiscovery is missing or true', async () => { const { fn, calls } = fakeFetch(); - const sink = createMeiliSink(() => BACKEND, fn); + const sink = createMeiliSink( + () => BACKEND, + () => null, + fn + ); await sink.onRecords( [ @@ -155,7 +194,11 @@ describe('createMeiliSink onRecords', () => { it('ignores records from other collections', async () => { const { fn, calls } = fakeFetch(); - const sink = createMeiliSink(() => BACKEND, fn); + const sink = createMeiliSink( + () => BACKEND, + () => null, + fn + ); await sink.onRecords( [ @@ -178,7 +221,11 @@ describe('createMeiliSink onRecords', () => { it('no-ops (no fetch) when the backend is unconfigured', async () => { const { fn, calls } = fakeFetch(); - const sink = createMeiliSink(() => null, fn); + const sink = createMeiliSink( + () => null, + () => null, + fn + ); await sink.onRecords([created('at://did:plc:alice/community.lexicon.calendar.event/3', {})], { phase: 'backfill' @@ -189,7 +236,7 @@ describe('createMeiliSink onRecords', () => { it('applies index settings once, before the first write (fresh-index safety)', async () => { const { fn, calls } = fakeFetch(); - const sink = createMeiliSink(() => BACKEND, fn); + const sink = createMeiliSink(() => BACKEND, () => null, fn); // Two batches on the same sink: a fresh-rollout `pnpm backfill` must not // let PUT /documents auto-create a bare index whose _geo/startsAt searches @@ -231,7 +278,11 @@ describe('fetch is invoked detached (workerd Illegal invocation guard)', () => { }); it('onRecords upsert/delete do not trip Illegal invocation', async () => { - const sink = createMeiliSink(() => BACKEND, strictFetch()); + const sink = createMeiliSink( + () => BACKEND, + () => null, + strictFetch() + ); await expect( sink.onRecords( [created('at://did:plc:alice/community.lexicon.calendar.event/4', { name: 'x' })], @@ -240,3 +291,144 @@ describe('fetch is invoked detached (workerd Illegal invocation guard)', () => { ).resolves.toBeUndefined(); }); }); + +/** A minimal D1 double: prepare().bind().all() returns the seeded resolved + * rows whose address_norm is in the bound args. Throw mode exercises the + * best-effort swallow. */ +function fakeDb( + rows: { address_norm: string; lat: number; lng: number }[], + mode: 'ok' | 'throw' = 'ok' +) { + return { + prepare() { + return { + bind(...args: unknown[]) { + return { + async all() { + if (mode === 'throw') throw new Error('no such table: geocode_cache'); + return { + results: rows.filter((r) => args.includes(r.address_norm)) as unknown as T[] + }; + } + }; + } + }; + } + } as unknown as D1Database; +} + +const ADDR_RECORD = { + name: 'TGIF meetup', + locations: [{ $type: ADDRESS_TYPE, name: 'TGIF meetup', locality: 'Bruxelles', country: 'BE' }] +}; + +describe('createMeiliSink geocode cache lookup', () => { + it('fills _geo from a resolved cache row for an address-only event', async () => { + const norm = normalizeAddress({ name: 'TGIF meetup', locality: 'Bruxelles', country: 'BE' })!; + const { fn, calls } = fakeFetch(); + const sink = createMeiliSink( + () => BACKEND, + () => fakeDb([{ address_norm: norm, lat: 50.84, lng: 4.36 }]), + fn + ); + + await sink.onRecords( + [created('at://did:plc:alice/community.lexicon.calendar.event/addr', ADDR_RECORD)], + { phase: 'live' } + ); + + const put = calls.find((c) => c.method === 'PUT'); + const docs = put!.body as Array>; + expect(docs[0]._geo).toEqual({ lat: 50.84, lng: 4.36 }); + }); + + it('leaves _geo unset on a cache miss but still indexes the doc', async () => { + const { fn, calls } = fakeFetch(); + const sink = createMeiliSink( + () => BACKEND, + () => fakeDb([]), + fn + ); + await sink.onRecords( + [created('at://did:plc:alice/community.lexicon.calendar.event/miss', ADDR_RECORD)], + { phase: 'live' } + ); + const docs = calls.find((c) => c.method === 'PUT')!.body as Array>; + expect(docs[0]._geo).toBeUndefined(); + expect(docs[0].name).toBe('TGIF meetup'); + }); + + it('does not consult the cache when the event already has coordinate _geo', async () => { + const norm = normalizeAddress({ locality: 'Bruxelles', country: 'BE' })!; + const { fn, calls } = fakeFetch(); + const sink = createMeiliSink( + () => BACKEND, + () => fakeDb([{ address_norm: norm, lat: 1, lng: 1 }]), + fn + ); + await sink.onRecords( + [ + created('at://did:plc:alice/community.lexicon.calendar.event/geo', { + locations: [ + { $type: 'community.lexicon.location.geo', latitude: '40.0', longitude: '-105.0' }, + { $type: ADDRESS_TYPE, locality: 'Bruxelles', country: 'BE' } + ] + }) + ], + { phase: 'live' } + ); + const docs = calls.find((c) => c.method === 'PUT')!.body as Array>; + expect(docs[0]._geo).toEqual({ lat: 40, lng: -105 }); // from .geo, not the cache row + }); + + it('swallows a cache query error and indexes without _geo', async () => { + const warn = vi.spyOn(console, 'warn').mockImplementation(() => {}); + const { fn, calls } = fakeFetch(); + const sink = createMeiliSink( + () => BACKEND, + () => fakeDb([], 'throw'), + fn + ); + await expect( + sink.onRecords( + [created('at://did:plc:alice/community.lexicon.calendar.event/err', ADDR_RECORD)], + { phase: 'live' } + ) + ).resolves.toBeUndefined(); + const docs = calls.find((c) => c.method === 'PUT')!.body as Array>; + expect(docs[0]._geo).toBeUndefined(); + warn.mockRestore(); + }); +}); + +describe('MeiliEventIndex.upsert', () => { + it('PUTs a full doc with _geo (the write the geocode job makes to attach coords)', async () => { + const { fn, calls } = fakeFetch(); + const doc = { + id: 'abc', + uri: 'at://did:plc:alice/community.lexicon.calendar.event/addr', + did: 'did:plc:alice', + rkey: 'addr', + name: 'Antwerp Drinks', + startsAt: '2099-01-01T10:00:00Z', + locationTypes: ['community.lexicon.location.address'], + _geo: { lat: 51.2194, lng: 4.4025 } + }; + await new MeiliEventIndex(BACKEND, fn).upsert([doc]); + + // PUT /documents is "add or update" — it merges an existing doc and lands a + // COMPLETE doc when absent, so a full-doc payload never leaves a {id,_geo} stub. + const put = calls.find( + (c) => + c.url === 'http://meili.local/indexes/events/documents?primaryKey=id' && c.method === 'PUT' + ); + expect(put).toBeDefined(); + expect(put!.body).toEqual([doc]); + }); + + it('no-ops on an empty doc list', async () => { + const { fn, calls } = fakeFetch(); + await new MeiliEventIndex(BACKEND, fn).upsert([]); + expect(calls).toHaveLength(0); + }); +}); diff --git a/apps/web/src/lib/search/server/meili-sink.ts b/apps/web/src/lib/search/server/meili-sink.ts index 18c1f4a..ca2219b 100644 --- a/apps/web/src/lib/search/server/meili-sink.ts +++ b/apps/web/src/lib/search/server/meili-sink.ts @@ -14,6 +14,8 @@ // so a Meilisearch outage degrades to "search index falls behind", not // "ingest stops". import { eventToSearchDoc, searchDocId, type SearchDoc } from './normalize'; +import { addressLocation, normalizeAddress } from './address-norm'; +import { isHiddenFromDiscovery } from './discoverability'; import type { ContrailConfig } from '@atmo-dev/contrail'; // The umbrella re-exports ContrailConfig (which carries `sinks?: Sink[]`) but @@ -25,14 +27,6 @@ type RecordEvent = Parameters[0][number]; /** The one collection we index for search. */ export const EVENT_COLLECTION = 'community.lexicon.calendar.event'; -/** True when the event author hid it from discovery. Mirrors the D1 filter in - * contrail.config.ts: only `preferences.showInDiscovery === false` hides; - * a missing/null field stays discoverable. */ -function isHiddenFromDiscovery(record: Record): boolean { - const prefs = record?.preferences as { showInDiscovery?: boolean } | undefined; - return prefs?.showInDiscovery === false; -} - export interface MeiliSinkBackend { url: string; apiKey: string; @@ -147,6 +141,34 @@ export async function applyMeiliSettings( await new MeiliEventIndex(backend, fetchFn).applySettings(); } +/** Best-effort fill of doc._geo from the geocode_cache. Read-only; a missing + * table or D1 hiccup just means "no _geo this pass" (never an ingest failure). + * Its real job is to reproduce the _geo the external geocode job wrote, so a + * later live UPDATE (full-doc PUT) of an already-geocoded event doesn't drop it. */ +async function fillGeoFromCache( + db: D1Database | null, + pending: { doc: SearchDoc; norm: string }[] +): Promise { + if (!db || pending.length === 0) return; + try { + const norms = [...new Set(pending.map((p) => p.norm))]; + const placeholders = norms.map(() => '?').join(','); + const { results } = await db + .prepare( + `SELECT address_norm, lat, lng FROM geocode_cache WHERE lat IS NOT NULL AND address_norm IN (${placeholders})` + ) + .bind(...norms) + .all<{ address_norm: string; lat: number; lng: number }>(); + const byNorm = new Map(results.map((r) => [r.address_norm, r])); + for (const { doc, norm } of pending) { + const row = byNorm.get(norm); + if (row) doc._geo = { lat: row.lat, lng: row.lng }; + } + } catch (e) { + console.warn('[search-sink] geocode cache lookup failed; indexing without _geo:', e); + } +} + /** Builds the contrail Sink. The backend is resolved lazily per batch via * `getBackend` because a Cloudflare Worker only has env per invocation, while * `contrail` is constructed once at module load — the cron/xrpc handlers set @@ -154,6 +176,7 @@ export async function applyMeiliSettings( * the backend is null (unconfigured), onRecords is a no-op. */ export function createMeiliSink( getBackend: () => MeiliSinkBackend | null, + getDb: () => D1Database | null = () => null, fetchFn?: typeof fetch ): Sink { // Apply the read-path's filterable/sortable settings once, before the first @@ -169,28 +192,39 @@ export function createMeiliSink( const docs: SearchDoc[] = []; const deletes: string[] = []; + // Docs that have no coordinate _geo but do carry an address — candidates + // for a geocode_cache fill. + const pending: { doc: SearchDoc; norm: string }[] = []; for (const e of events) { if (e.collection !== EVENT_COLLECTION) continue; - // Mirror the D1 discoverable filter (contrail.config.ts): only - // `showInDiscovery === false` hides; missing/null stays discoverable. + // Discoverability is decided by the shared predicate (discoverability.ts) + // so this in-memory filter and the D1 SQL filter never diverge. // Index discoverable creates; for everything else — real deletes AND // events hidden from discovery — remove the doc so the search index // never holds a hidden event's name/description and a discoverable→ // unlisted flip purges the existing entry. if (e.kind === 'created' && !isHiddenFromDiscovery(e.record)) { - docs.push( - eventToSearchDoc({ - uri: e.uri, - did: e.did, - collection: e.collection, - rkey: e.rkey, - record: e.record - }) - ); + const doc = eventToSearchDoc({ + uri: e.uri, + did: e.did, + collection: e.collection, + rkey: e.rkey, + record: e.record + }); + docs.push(doc); + if (!doc._geo) { + const loc = addressLocation(e.record); + const norm = loc ? normalizeAddress(loc) : null; + if (norm) pending.push({ doc, norm }); + } } else { deletes.push(searchDocId(e.uri)); } } + + // Mutates doc._geo in place before the upsert below. + await fillGeoFromCache(getDb(), pending); + if (docs.length === 0 && deletes.length === 0) return; const index = new MeiliEventIndex(backend, fetchFn); diff --git a/apps/web/src/lib/search/server/normalize.test.ts b/apps/web/src/lib/search/server/normalize.test.ts index 82387ff..1926467 100644 --- a/apps/web/src/lib/search/server/normalize.test.ts +++ b/apps/web/src/lib/search/server/normalize.test.ts @@ -1,10 +1,16 @@ import { describe, it, expect } from 'vitest'; -import { eventToSearchDoc, type EventRecordPayload } from './normalize'; +import { eventToSearchDoc, recordGeo, type EventRecordPayload } from './normalize'; const URI = 'at://did:plc:alice/community.lexicon.calendar.event/1'; function payload(record: Record): EventRecordPayload { - return { uri: URI, did: 'did:plc:alice', collection: 'community.lexicon.calendar.event', rkey: '1', record }; + return { + uri: URI, + did: 'did:plc:alice', + collection: 'community.lexicon.calendar.event', + rkey: '1', + record + }; } function geoLoc(latitude: string, longitude: string) { @@ -45,3 +51,38 @@ describe('eventToSearchDoc geo derivation', () => { expect(doc._geo).toBeUndefined(); }); }); + +const FSQ = 'community.lexicon.location.fsq'; +const ADDRESS = 'community.lexicon.location.address'; + +describe('recordGeo', () => { + it('derives coordinates from a geo location', () => { + expect(recordGeo({ locations: [geoLoc('40.0', '-105.0')] })).toEqual({ lat: 40, lng: -105 }); + }); + + it('derives coordinates from an fsq location that carries lat/lng', () => { + expect( + recordGeo({ locations: [{ $type: FSQ, latitude: '50.84', longitude: '4.36' }] }) + ).toEqual({ lat: 50.84, lng: 4.36 }); + }); + + it('returns undefined when the only coordinate location is out of range', () => { + expect(recordGeo({ locations: [geoLoc('999', '-105.0')] })).toBeUndefined(); + }); + + it('returns undefined for an address-only record', () => { + expect( + recordGeo({ locations: [{ $type: ADDRESS, locality: 'Dayton', country: 'US' }] }) + ).toBeUndefined(); + }); + + it('matches the _geo eventToSearchDoc derives (single source of truth)', () => { + const record = { + locations: [ + { $type: FSQ, latitude: '50.84', longitude: '4.36' }, + { $type: ADDRESS, locality: 'x' } + ] + }; + expect(recordGeo(record)).toEqual(eventToSearchDoc(payload(record))._geo); + }); +}); diff --git a/apps/web/src/lib/search/server/normalize.ts b/apps/web/src/lib/search/server/normalize.ts index 1a2a65a..ab9e610 100644 --- a/apps/web/src/lib/search/server/normalize.ts +++ b/apps/web/src/lib/search/server/normalize.ts @@ -97,11 +97,29 @@ function str(v: unknown): string | undefined { return typeof v === 'string' ? v : undefined; } -export function eventToSearchDoc(payload: EventRecordPayload): SearchDoc { - const record = payload.record ?? {}; - const locations = (Array.isArray(record.locations) ? record.locations : []).filter( +/** The record's locations[] narrowed to the location objects deriveGeo and the + * doc builder read. */ +function recordLocations(record: Record): Loc[] { + return (Array.isArray(record?.locations) ? record.locations : []).filter( (l): l is Loc => !!l && typeof l === 'object' ); +} + +/** The single canonical _geo a record resolves to (precedence geo > fsq > + * hthree, in-range only), or undefined. Shared by the search doc AND the + * external geocode worklist so "already has coordinates" means the exact same + * thing in both: the worklist must not geocode an event the index already + * geo-derives (which would overwrite precise fsq/geo coords), and must still + * geocode one whose only coordinate location is out of range (no _geo). */ +export function recordGeo( + record: Record +): { lat: number; lng: number } | undefined { + return deriveGeo(recordLocations(record)); +} + +export function eventToSearchDoc(payload: EventRecordPayload): SearchDoc { + const record = payload.record ?? {}; + const locations = recordLocations(record); const doc: SearchDoc = { id: searchDocId(payload.uri),