packages/ai/components/api/pineconeUpload.ts
1'use server'
2
3import { createHash, timingSafeEqual } from 'node:crypto'
4import type { Index, RecordMetadata } from '@pinecone-database/pinecone'
5import { Pinecone } from '@pinecone-database/pinecone'
6import { extractText } from 'unpdf'
7import { canonicalAgendaUrl } from '../agendaUrl'
8import {
9 CHUNK_OVERLAP,
10 CHUNK_SIZE,
11 chunkAgendaMarkdown,
12 chunkText,
13} from './chunking'
14
15// Polyfill: unpdf 1.6 (pdf.js) requires `Promise.try` (V8 13.x / Node 23+).
16// On Node 22 this throws inside PDF parsing and crashes the dev worker.
17if (typeof (Promise as unknown as { try?: unknown }).try !== 'function') {
18 ;(
19 Promise as unknown as {
20 try: (
21 fn: (...args: unknown[]) => unknown,
22 ...args: unknown[]
23 ) => Promise<unknown>
24 }
25 ).try = <T>(
26 fn: (...args: unknown[]) => T | PromiseLike<T>,
27 ...args: unknown[]
28 ): Promise<T> => new Promise<T>((resolve) => resolve(fn(...args)))
29}
30
31const LOG_PREFIX = '[Pinecone]'
32
33function log(msg: string, data?: Record<string, unknown>) {
34 const payload = data ? ` ${JSON.stringify(data)}` : ''
35 console.log(`${LOG_PREFIX} ${msg}${payload}`)
36}
37
38function logError(msg: string, error: unknown) {
39 const err = error instanceof Error ? error : new Error(String(error))
40 console.error(`${LOG_PREFIX} ${msg}`, err.message, err.stack ?? '')
41}
42
43type ScopeType = 'general' | 'event'
44
45type ChunkRecord = {
46 id: string
47 fileName: string
48 chunkIndex: number
49 text: string
50}
51
52const DEFAULTS = {
53 batchSize: 50,
54 embeddingBatchSize: 100,
55 chunkSize: CHUNK_SIZE,
56 chunkOverlap: CHUNK_OVERLAP,
57 embeddingModel: 'text-embedding-3-small',
58 dimension: 1536,
59 retries: 4,
60 retryBaseDelayMs: 1000,
61 /** query() topK ceiling used when finding a file's existing chunk ids */
62 maxChunksPerFile: 2000,
63} as const
64
65function requireEnv(name: string): string {
66 const value = process.env[name]?.trim()
67 if (!value) {
68 throw new Error(`Missing environment variable: ${name}`)
69 }
70 return value
71}
72
73/**
74 * Client + index handles are cached per process: chat tool calls hit Pinecone
75 * on every search, and rebuilding the client per call costs connection setup
76 * on the hot path.
77 */
78let pineconeClient: Pinecone | undefined
79const indexHandles = new Map<string, Index>()
80
81function getPinecone(): Pinecone {
82 if (!pineconeClient) {
83 pineconeClient = new Pinecone({ apiKey: requireEnv('PINECONE_API_KEY') })
84 }
85 return pineconeClient
86}
87
88function getIndex(indexName: string): Index {
89 let handle = indexHandles.get(indexName)
90 if (!handle) {
91 handle = getPinecone().index(indexName)
92 indexHandles.set(indexName, handle)
93 }
94 return handle
95}
96
97/**
98 * Namespace consolidation switch (SoW §8.2): when PINECONE_CONSOLIDATED_INDEX
99 * is set, every store lives as a namespace (named after its legacy index)
100 * inside that single index. Until the env var is set — after running
101 * scripts/pinecone-migrate-namespaces — behaviour is identical to today.
102 */
103function getTargetIndex(indexName: string): Index {
104 const consolidated = process.env.PINECONE_CONSOLIDATED_INDEX?.trim()
105 if (!consolidated) return getIndex(indexName)
106 const key = `${consolidated}::${indexName}`
107 let handle = indexHandles.get(key)
108 if (!handle) {
109 handle = getPinecone().index(consolidated).namespace(indexName)
110 indexHandles.set(key, handle)
111 }
112 return handle
113}
114
115/**
116 * Shared-secret auth for mutation endpoints. Requests require a configured
117 * INTERNAL_API_SECRET and a matching Authorization bearer token.
118 * Returns an error Response, or null when the request may proceed.
119 */
120export function checkPineconeAuth(request: Request): Response | null {
121 const secret = process.env.INTERNAL_API_SECRET?.trim()
122 if (!secret) {
123 console.warn(
124 `${LOG_PREFIX} INTERNAL_API_SECRET is not set; mutation endpoint is disabled`,
125 )
126 return Response.json(
127 { ok: false, error: 'Mutation authentication is not configured' },
128 { status: 503 },
129 )
130 }
131 const header = request.headers.get('authorization') ?? ''
132 const expected = Buffer.from(`Bearer ${secret}`)
133 const received = Buffer.from(header)
134 if (
135 expected.length === received.length &&
136 timingSafeEqual(expected, received)
137 ) {
138 return null
139 }
140 log('Auth rejected: missing/invalid bearer token')
141 return Response.json({ ok: false, error: 'Unauthorized' }, { status: 401 })
142}
143
144/** Only network/429/5xx failures are worth retrying; 4xx are permanent. */
145function isTransientError(error: unknown): boolean {
146 if (!(error instanceof Error)) return true
147 if (
148 /PineconeConnection|ECONNRESET|ETIMEDOUT|fetch failed|network/i.test(
149 `${error.name} ${error.message}`,
150 )
151 ) {
152 return true
153 }
154 const statusMatch = error.message.match(/\((\d{3})\)/)
155 if (statusMatch) {
156 const status = Number(statusMatch[1])
157 return status === 429 || status >= 500
158 }
159 // SDK errors that mention explicit client-side rejection are permanent.
160 return !/invalid|unauthorized|forbidden|not found|400|401|403|404/i.test(
161 error.message,
162 )
163}
164
165/** Retry transient failures (network, 429, 5xx) with exponential backoff. */
166async function withRetry<T>(label: string, fn: () => Promise<T>): Promise<T> {
167 let lastError: unknown
168 for (let attempt = 0; attempt < DEFAULTS.retries; attempt += 1) {
169 try {
170 return await fn()
171 } catch (error) {
172 lastError = error
173 const isLast = attempt === DEFAULTS.retries - 1
174 if (isLast || !isTransientError(error)) break
175 const delay = DEFAULTS.retryBaseDelayMs * 2 ** attempt
176 log(`${label} failed, retrying`, { attempt: attempt + 1, delayMs: delay })
177 await new Promise((resolve) => setTimeout(resolve, delay))
178 }
179 }
180 throw lastError
181}
182
183function sliceIntoBatches<T>(items: T[], batchSize: number): T[][] {
184 const batches: T[][] = []
185 for (let i = 0; i < items.length; i += batchSize) {
186 batches.push(items.slice(i, i + batchSize))
187 }
188 return batches
189}
190
191function createChunkId(
192 scopeKey: string,
193 fileName: string,
194 chunkIndex: number,
195): string {
196 return createHash('sha1')
197 .update(`${scopeKey}::${fileName}::${chunkIndex}`)
198 .digest('hex')
199 .slice(0, 32)
200}
201
202function isSupportedFile(file: File): boolean {
203 const name = file.name.toLowerCase()
204 return name.endsWith('.pdf') || name.endsWith('.txt')
205}
206
207async function readTextFromFile(file: File): Promise<string | null> {
208 const lowerName = file.name.toLowerCase()
209 if (lowerName.endsWith('.pdf')) {
210 const bytes = new Uint8Array(await file.arrayBuffer())
211 const { text } = await extractText(bytes, { mergePages: true })
212 return text.trim() || null
213 }
214
215 if (lowerName.endsWith('.txt')) {
216 const text = (await file.text()).trim()
217 return text || null
218 }
219
220 return null
221}
222
223async function embedTexts(
224 texts: string[],
225 openAiApiKey: string,
226 embeddingModel: string,
227): Promise<number[][]> {
228 const json = await withRetry('Embedding request', async () => {
229 const response = await fetch('https://api.openai.com/v1/embeddings', {
230 method: 'POST',
231 headers: {
232 'Content-Type': 'application/json',
233 Authorization: `Bearer ${openAiApiKey}`,
234 },
235 body: JSON.stringify({
236 model: embeddingModel,
237 input: texts,
238 }),
239 })
240
241 if (!response.ok) {
242 const errorText = await response.text()
243 throw new Error(
244 `Embedding request failed (${response.status}): ${errorText.slice(0, 800)}`,
245 )
246 }
247
248 return (await response.json()) as {
249 data?: Array<{ embedding: number[] }>
250 }
251 })
252
253 const embeddings = json.data?.map((item) => item.embedding) || []
254 if (embeddings.length !== texts.length) {
255 throw new Error(
256 `Embedding response size mismatch. Expected ${texts.length}, received ${embeddings.length}`,
257 )
258 }
259
260 for (const embedding of embeddings) {
261 if (embedding.length !== DEFAULTS.dimension) {
262 throw new Error(
263 `Embedding dimension mismatch. Expected ${DEFAULTS.dimension}, received ${embedding.length}.`,
264 )
265 }
266 }
267
268 return embeddings
269}
270
271type PineconeVector = {
272 id: string
273 values: number[]
274 metadata: RecordMetadata
275}
276
277async function upsertVectors(index: Index, vectors: PineconeVector[]) {
278 let uploadedCount = 0
279 for (const batch of sliceIntoBatches(vectors, DEFAULTS.batchSize)) {
280 await withRetry('Pinecone upsert', () => index.upsert({ records: batch }))
281 uploadedCount += batch.length
282 }
283 return uploadedCount
284}
285
286/** List vector ids in the index's namespace, optionally scoped to an id prefix. */
287async function listAllVectorIds(
288 index: Index,
289 idPrefix?: string,
290): Promise<Set<string>> {
291 const ids = new Set<string>()
292 let paginationToken: string | undefined
293 do {
294 const page = await withRetry('Pinecone list', () =>
295 index.listPaginated({ limit: 100, paginationToken, prefix: idPrefix }),
296 )
297 for (const vector of page.vectors ?? []) {
298 if (vector.id) ids.add(vector.id)
299 }
300 paginationToken = page.pagination?.next
301 } while (paginationToken)
302 return ids
303}
304
305/**
306 * Delete every id in the index that is not part of the freshly upserted set.
307 * This is what makes refreshes true replacements: without it, shrinking
308 * content leaves stale tail vectors behind forever.
309 */
310async function deleteStaleVectors(
311 index: Index,
312 freshIds: Set<string>,
313 idPrefix?: string,
314): Promise<number> {
315 const existing = await listAllVectorIds(index, idPrefix)
316 const stale = [...existing].filter((id) => !freshIds.has(id))
317 if (stale.length === 0) return 0
318 for (const batch of sliceIntoBatches(stale, 1000)) {
319 // v7 SDK takes an options object; a bare array fails with "Invalid request".
320 await withRetry('Pinecone delete stale', () =>
321 index.deleteMany({ ids: batch }),
322 )
323 }
324 return stale.length
325}
326
327/**
328 * Find existing chunk ids for a file via a metadata-filtered query (serverless
329 * indexes do not support metadata-filter deletes, but filtered queries work).
330 */
331async function findFileVectorIds(
332 index: Index,
333 scopeKey: string,
334 fileName: string,
335): Promise<string[]> {
336 // Cosine indexes reject all-zero query vectors; any unit vector works since
337 // the metadata filter does the selection and scores are irrelevant here.
338 const probeVector = new Array(DEFAULTS.dimension).fill(0)
339 probeVector[0] = 1
340 const response = await withRetry('Pinecone file query', () =>
341 index.query({
342 vector: probeVector,
343 topK: DEFAULTS.maxChunksPerFile,
344 filter: { scopeKey: { $eq: scopeKey }, fileName: { $eq: fileName } },
345 includeValues: false,
346 includeMetadata: false,
347 }),
348 )
349 return (response.matches ?? []).map((match) => match.id)
350}
351
352export async function handlePineconeUpload(
353 request: Request,
354): Promise<Response> {
355 try {
356 const form = await request.formData()
357 const scopeType = (form.get('scopeType') as ScopeType) || 'general'
358 const scopeKey = String(form.get('scopeKey') || 'general').trim() || 'general'
359 const eventDocumentId = String(form.get('eventDocumentId') || '').trim()
360 const eventTitle = String(form.get('eventTitle') || '').trim()
361 const eventSlug = String(form.get('eventSlug') || '').trim()
362 const eventInwinkId = String(form.get('eventInwinkId') || '').trim()
363
364 log('Upload started', { scopeKey, scopeType })
365
366 if (scopeType !== 'general' && scopeType !== 'event') {
367 log('Upload rejected: invalid scopeType', { scopeType })
368 return Response.json({ error: 'Invalid scopeType' }, { status: 400 })
369 }
370
371 if (
372 !VALID_SCOPE_KEYS.includes(scopeKey as (typeof VALID_SCOPE_KEYS)[number])
373 ) {
374 log('Upload rejected: invalid scopeKey', { scopeKey })
375 return Response.json(
376 {
377 error: `Invalid scopeKey. Must be one of: ${VALID_SCOPE_KEYS.join(', ')}`,
378 },
379 { status: 400 },
380 )
381 }
382
383 const files = form
384 .getAll('files')
385 .filter((entry): entry is File => entry instanceof File && entry.size > 0)
386
387 if (files.length === 0) {
388 log('Upload rejected: no files')
389 return Response.json({ error: 'No files received' }, { status: 400 })
390 }
391
392 const supportedFiles = files.filter(isSupportedFile)
393 const skippedFiles = files
394 .filter((file) => !isSupportedFile(file))
395 .map((file) => file.name)
396 if (supportedFiles.length === 0) {
397 log('Upload rejected: no supported files', { skipped: skippedFiles })
398 return Response.json(
399 { error: 'No supported files. Only .pdf and .txt are allowed.' },
400 { status: 400 },
401 )
402 }
403
404 const openAiApiKey = requireEnv('OPENAI_API_KEY')
405 const indexName = scopeKeyToIndexName(scopeKey)
406 const index = getTargetIndex(indexName)
407 log('Upload target index (default namespace)', { scopeKey, indexName })
408
409 // Same-name uploads are rejected unless explicitly overwriting, so a
410 // second editor can't silently replace someone else's live document
411 // (versioning is owned by the client per the SoW).
412 const overwrite = String(form.get('overwrite') || '').trim() === 'true'
413 if (!overwrite) {
414 const duplicates: string[] = []
415 for (const file of supportedFiles) {
416 try {
417 const existing = await findFileVectorIds(index, scopeKey, file.name)
418 if (existing.length > 0) duplicates.push(file.name)
419 } catch (error) {
420 logError(`Duplicate check failed for ${file.name}`, error)
421 }
422 }
423 if (duplicates.length > 0) {
424 log('Upload rejected: duplicate file names', { scopeKey, duplicates })
425 // `duplicates` lets the Studio offer a confirm-and-replace retry
426 // (re-upload with overwrite=true) instead of a dead end.
427 return Response.json(
428 {
429 error: `File(s) already in this store: ${duplicates.join(', ')}. Upload again with overwrite to replace them.`,
430 duplicates,
431 },
432 { status: 409 },
433 )
434 }
435 }
436
437 const chunks: ChunkRecord[] = []
438 const freshIdsPerFile = new Map<string, Set<string>>()
439 let processedFiles = 0
440
441 for (const file of supportedFiles) {
442 const text = await readTextFromFile(file)
443 if (!text) continue
444 processedFiles += 1
445 const textChunks = chunkText(text, DEFAULTS.chunkSize, DEFAULTS.chunkOverlap)
446 const fileIds = new Set<string>()
447 for (let i = 0; i < textChunks.length; i += 1) {
448 const id = createChunkId(scopeKey, file.name, i)
449 fileIds.add(id)
450 chunks.push({
451 id,
452 fileName: file.name,
453 chunkIndex: i,
454 text: textChunks[i],
455 })
456 }
457 freshIdsPerFile.set(file.name, fileIds)
458 }
459
460 if (chunks.length === 0) {
461 log('Upload completed: no text extracted', { scopeKey, processedFiles })
462 return Response.json({
463 scopeType,
464 scopeKey,
465 processedFiles,
466 totalVectors: 0,
467 skippedFiles,
468 })
469 }
470
471 const embeddingBatches = sliceIntoBatches(chunks, DEFAULTS.embeddingBatchSize)
472 let uploadedCount = 0
473
474 for (const chunkBatch of embeddingBatches) {
475 const embeddings = await embedTexts(
476 chunkBatch.map((chunk) => chunk.text),
477 openAiApiKey,
478 DEFAULTS.embeddingModel,
479 )
480
481 const vectors = chunkBatch.map((chunk, idx) => {
482 const metadata: RecordMetadata = {
483 scopeType,
484 scopeKey,
485 fileName: chunk.fileName,
486 chunkIndex: chunk.chunkIndex,
487 text: chunk.text,
488 }
489 if (eventDocumentId) metadata.eventDocumentId = eventDocumentId
490 if (eventTitle) metadata.eventTitle = eventTitle
491 if (eventSlug) metadata.eventSlug = eventSlug
492 if (eventInwinkId) metadata.eventInwinkId = eventInwinkId
493 return { id: chunk.id, values: embeddings[idx], metadata }
494 })
495
496 uploadedCount += await upsertVectors(index, vectors)
497 }
498
499 // Re-uploading a shorter version of a file must not leave stale tail
500 // chunks behind. Best-effort: an error here doesn't fail the upload.
501 let staleDeleted = 0
502 for (const [fileName, freshIds] of freshIdsPerFile) {
503 try {
504 const existingIds = await findFileVectorIds(index, scopeKey, fileName)
505 const stale = existingIds.filter((id) => !freshIds.has(id))
506 if (stale.length > 0) {
507 await withRetry('Pinecone delete stale file chunks', () =>
508 index.deleteMany({ ids: stale }),
509 )
510 staleDeleted += stale.length
511 }
512 } catch (error) {
513 logError(`Stale chunk cleanup failed for ${fileName}`, error)
514 }
515 }
516
517 log('Upload success', {
518 scopeKey,
519 indexName,
520 processedFiles,
521 totalVectors: uploadedCount,
522 staleDeleted,
523 skippedCount: skippedFiles.length,
524 })
525 return Response.json({
526 scopeType,
527 scopeKey,
528 processedFiles,
529 totalVectors: uploadedCount,
530 skippedFiles,
531 })
532 } catch (error) {
533 logError('Upload failed', error)
534 const message = error instanceof Error ? error.message : 'Unexpected error'
535 return Response.json({ error: message }, { status: 500 })
536 }
537}
538
539const VALID_SCOPE_KEYS = [
540 'general',
541 'event-world',
542 'event-america',
543 'sales-agent',
544 'sales-agent-spex',
545] as const
546
547const AGENDA_SCOPE_KEYS = ['event-world', 'event-america'] as const
548
549/** Index names matching your Pinecone hosts (default namespace only, no custom namespace). */
550const INDEX_NAMES = {
551 general: 'unleash-faq',
552 'event-world': 'event-world',
553 'event-america': 'event-america',
554 'event-world-agenda': 'event-world-agenda',
555 'event-america-agenda': 'event-america-agenda',
556 sitemap: 'unleash-sitemap',
557 media: 'media-catalogue',
558 'sales-agent': 'sales-agent',
559 'sales-agent-spex': 'sales-agent-spex',
560} as const
561
562function getDefaultIndexName(): string {
563 return process.env.PINECONE_INDEX_NAME?.trim() || INDEX_NAMES.general
564}
565
566/**
567 * Map scopeKey to Pinecone index name for file uploads and "Clear store".
568 * Uses default namespace only (no namespace parameter sent).
569 */
570export function scopeKeyToIndexName(scopeKey: string): string {
571 switch (scopeKey) {
572 case 'general':
573 return process.env.PINECONE_INDEX_GENERAL?.trim() || getDefaultIndexName()
574 case 'event-world':
575 return (
576 process.env.PINECONE_INDEX_EVENT_WORLD?.trim() || INDEX_NAMES['event-world']
577 )
578 case 'event-america':
579 return (
580 process.env.PINECONE_INDEX_EVENT_AMERICA?.trim() ||
581 INDEX_NAMES['event-america']
582 )
583 case 'sales-agent':
584 return (
585 process.env.PINECONE_INDEX_SALES_AGENT?.trim() || INDEX_NAMES['sales-agent']
586 )
587 case 'sales-agent-spex':
588 return (
589 process.env.PINECONE_INDEX_SALES_AGENT_SPEX?.trim() ||
590 INDEX_NAMES['sales-agent-spex']
591 )
592 default:
593 return getDefaultIndexName()
594 }
595}
596
597/**
598 * Map scopeKey to Pinecone index name for agenda (Refresh World/America Agenda).
599 * Uses default namespace only.
600 */
601export function scopeKeyToAgendaIndexName(
602 scopeKey: 'event-world' | 'event-america',
603): string {
604 switch (scopeKey) {
605 case 'event-world':
606 return (
607 process.env.PINECONE_INDEX_EVENT_WORLD_AGENDA?.trim() ||
608 INDEX_NAMES['event-world-agenda']
609 )
610 case 'event-america':
611 return (
612 process.env.PINECONE_INDEX_EVENT_AMERICA_AGENDA?.trim() ||
613 INDEX_NAMES['event-america-agenda']
614 )
615 default:
616 return INDEX_NAMES['event-world-agenda']
617 }
618}
619
620export function sitemapIndexName(): string {
621 return process.env.PINECONE_INDEX_SITEMAP?.trim() || INDEX_NAMES.sitemap
622}
623
624export function mediaIndexName(): string {
625 return process.env.PINECONE_INDEX_MEDIA?.trim() || INDEX_NAMES.media
626}
627
628export type PineconeStoreStatus = {
629 store: string
630 indexName: string
631 vectorCount: number | null
632 error?: string
633}
634
635/**
636 * Resolved index name + live vector count for every store the Studio manages.
637 * Read-only; used by the Studio tool and for health checks.
638 */
639export async function getPineconeStoreStatuses(): Promise<
640 PineconeStoreStatus[]
641> {
642 const pinecone = getPinecone()
643
644 const stores: Array<{ store: string; indexName: string }> = [
645 ...VALID_SCOPE_KEYS.map((key) => ({
646 store: key,
647 indexName: scopeKeyToIndexName(key),
648 })),
649 ...AGENDA_SCOPE_KEYS.map((key) => ({
650 store: `${key}-agenda`,
651 indexName: scopeKeyToAgendaIndexName(key),
652 })),
653 { store: 'media', indexName: mediaIndexName() },
654 { store: 'sitemap', indexName: sitemapIndexName() },
655 ]
656
657 // Consolidated mode: one stats call, counts read per namespace.
658 const consolidated = process.env.PINECONE_CONSOLIDATED_INDEX?.trim()
659 if (consolidated) {
660 try {
661 const stats = await withRetry('Pinecone stats', () =>
662 pinecone.index(consolidated).describeIndexStats(),
663 )
664 return stores.map(({ store, indexName }) => ({
665 store,
666 indexName: `${consolidated}/${indexName}`,
667 vectorCount: stats.namespaces?.[indexName]?.recordCount ?? 0,
668 }))
669 } catch (error) {
670 const message = error instanceof Error ? error.message : String(error)
671 return stores.map(({ store, indexName }) => ({
672 store,
673 indexName: `${consolidated}/${indexName}`,
674 vectorCount: null,
675 error: message,
676 }))
677 }
678 }
679
680 return Promise.all(
681 stores.map(async ({ store, indexName }) => {
682 try {
683 const stats = await withRetry('Pinecone stats', () =>
684 pinecone.index(indexName).describeIndexStats(),
685 )
686 return {
687 store,
688 indexName,
689 vectorCount: stats.totalRecordCount ?? 0,
690 }
691 } catch (error) {
692 const message = error instanceof Error ? error.message : String(error)
693 return { store, indexName, vectorCount: null, error: message }
694 }
695 }),
696 )
697}
698
699export type PineconeClearResponse =
700 | { ok: true; scopeKey: string }
701 | { ok: false; error: string }
702
703/**
704 * Clears the default namespace of the index for the given scopeKey.
705 * @param forAgenda - when true, clears the agenda index (event-world-agenda / event-america-agenda); when false, clears the file index (event-world / event-america / general).
706 */
707export async function clearPineconeScope(
708 scopeKey: string,
709 options?: { forAgenda?: boolean },
710): Promise<void> {
711 const indexName =
712 options?.forAgenda &&
713 (scopeKey === 'event-world' || scopeKey === 'event-america')
714 ? scopeKeyToAgendaIndexName(scopeKey as 'event-world' | 'event-america')
715 : scopeKeyToIndexName(scopeKey)
716 log('Clear scope started', {
717 scopeKey,
718 indexName,
719 forAgenda: options?.forAgenda,
720 })
721 const index = getTargetIndex(indexName)
722 try {
723 await index.deleteAll()
724 } catch (error) {
725 // Pinecone serverless returns 404 when the namespace is already empty.
726 if (!(error instanceof Error && /404|not found/i.test(error.message))) {
727 throw error
728 }
729 log('Clear scope: namespace already empty', { scopeKey, indexName })
730 }
731 log('Clear scope done', { scopeKey, indexName })
732}
733
734export async function handlePineconeClear(request: Request): Promise<Response> {
735 try {
736 const form = await request.formData()
737 const scopeKey = String(form.get('scopeKey') || '').trim()
738
739 if (!scopeKey) {
740 log('Clear rejected: missing scopeKey')
741 return Response.json(
742 { ok: false, error: 'Missing scopeKey' } satisfies PineconeClearResponse,
743 { status: 400 },
744 )
745 }
746
747 if (
748 !VALID_SCOPE_KEYS.includes(scopeKey as (typeof VALID_SCOPE_KEYS)[number])
749 ) {
750 log('Clear rejected: invalid scopeKey', { scopeKey })
751 return Response.json(
752 {
753 ok: false,
754 error: `Invalid scopeKey. Must be one of: ${VALID_SCOPE_KEYS.join(', ')}`,
755 } satisfies PineconeClearResponse,
756 { status: 400 },
757 )
758 }
759
760 await clearPineconeScope(scopeKey)
761
762 log('Clear success', { scopeKey })
763 return Response.json({
764 ok: true,
765 scopeKey,
766 } satisfies PineconeClearResponse)
767 } catch (error) {
768 logError('Clear failed', error)
769 const message = error instanceof Error ? error.message : 'Unexpected error'
770 return Response.json(
771 { ok: false, error: message } satisfies PineconeClearResponse,
772 { status: 500 },
773 )
774 }
775}
776
777export type PineconeUpsertTextResponse =
778 | { ok: true; scopeKey: string; totalVectors: number }
779 | { ok: false; error: string }
780
781/**
782 * Chunks the given text, embeds it, and upserts into the agenda index for the
783 * given event scopeKey, then deletes any leftover vectors from previous longer
784 * agendas (the index is dedicated to agenda content, so a full sync is safe).
785 * Embedding happens before any deletion so a failure never empties the store.
786 */
787export async function upsertTextContent(
788 scopeKey: 'event-world' | 'event-america',
789 text: string,
790): Promise<{ totalVectors: number }> {
791 if (!AGENDA_SCOPE_KEYS.includes(scopeKey)) {
792 throw new Error(
793 'Invalid scopeKey for agenda. Use event-world or event-america.',
794 )
795 }
796
797 const indexName = scopeKeyToAgendaIndexName(scopeKey)
798 log('Upsert text started', { scopeKey, indexName, textLength: text.length })
799
800 const openAiApiKey = requireEnv('OPENAI_API_KEY')
801 const index = getTargetIndex(indexName)
802
803 // Session-aware chunking: every chunk is a complete session block, so
804 // retrieval never returns a session without its time/room/speaker lines.
805 const textChunks = chunkAgendaMarkdown(text)
806 if (textChunks.length === 0) {
807 log('Upsert text: no chunks', { scopeKey, indexName })
808 return { totalVectors: 0 }
809 }
810
811 log('Upsert text: chunked', {
812 scopeKey,
813 indexName,
814 chunkCount: textChunks.length,
815 })
816
817 const chunks = textChunks.map((t, i) => ({
818 id: createChunkId(scopeKey, 'agenda', i),
819 fileName: 'agenda',
820 chunkIndex: i,
821 text: t,
822 }))
823 const freshIds = new Set(chunks.map((chunk) => chunk.id))
824
825 let uploadedCount = 0
826 for (const chunkBatch of sliceIntoBatches(
827 chunks,
828 DEFAULTS.embeddingBatchSize,
829 )) {
830 const embeddings = await embedTexts(
831 chunkBatch.map((c) => c.text),
832 openAiApiKey,
833 DEFAULTS.embeddingModel,
834 )
835
836 const vectors = chunkBatch.map((chunk, idx) => ({
837 id: chunk.id,
838 values: embeddings[idx],
839 metadata: {
840 scopeType: 'event',
841 scopeKey,
842 fileName: chunk.fileName,
843 chunkIndex: chunk.chunkIndex,
844 text: chunk.text,
845 } satisfies RecordMetadata,
846 }))
847
848 uploadedCount += await upsertVectors(index, vectors)
849 }
850
851 const staleDeleted = await deleteStaleVectors(index, freshIds)
852 agendaCache.delete(scopeKey)
853
854 log('Upsert text success', {
855 scopeKey,
856 indexName,
857 totalVectors: uploadedCount,
858 staleDeleted,
859 })
860 return { totalVectors: uploadedCount }
861}
862
863export type CatalogueRecord = {
864 id: string
865 text: string
866 metadata: Record<string, string>
867}
868
869/**
870 * Full-sync a catalogue-style index (one vector per record, no chunking):
871 * embed + upsert all records, then delete anything no longer present.
872 * Embedding happens before any deletion so a failure never empties the store.
873 */
874async function syncCatalogueIndex(
875 indexName: string,
876 records: CatalogueRecord[],
877 options?: {
878 /** Scope stale cleanup to ids with this prefix — required when the
879 * records share an index/namespace with other content (e.g. FAQ
880 * entries living beside uploaded files in the general store). */
881 idPrefix?: string
882 },
883): Promise<{ totalVectors: number; staleDeleted: number }> {
884 if (records.length === 0) {
885 // Removing the last FAQ/page/media item must remove its old knowledge.
886 // Restrict cleanup to catalogue-owned IDs so uploaded files are preserved.
887 if (!options?.idPrefix)
888 throw new Error('Empty catalogue sync requires an ID prefix')
889 const staleDeleted = await deleteStaleVectors(
890 getTargetIndex(indexName),
891 new Set(),
892 options.idPrefix,
893 )
894 return { totalVectors: 0, staleDeleted }
895 }
896
897 log('Catalogue sync started', { indexName, recordCount: records.length })
898
899 const openAiApiKey = requireEnv('OPENAI_API_KEY')
900 const index = getTargetIndex(indexName)
901
902 let uploadedCount = 0
903 for (const batch of sliceIntoBatches(records, DEFAULTS.embeddingBatchSize)) {
904 const embeddings = await embedTexts(
905 batch.map((r) => r.text),
906 openAiApiKey,
907 DEFAULTS.embeddingModel,
908 )
909
910 const vectors = batch.map((record, idx) => ({
911 id: record.id,
912 values: embeddings[idx],
913 metadata: { ...record.metadata, text: record.text },
914 }))
915
916 uploadedCount += await upsertVectors(index, vectors)
917 }
918
919 const freshIds = new Set(records.map((record) => record.id))
920 const staleDeleted = await deleteStaleVectors(
921 index,
922 freshIds,
923 options?.idPrefix,
924 )
925
926 log('Catalogue sync success', {
927 indexName,
928 totalVectors: uploadedCount,
929 staleDeleted,
930 })
931 return { totalVectors: uploadedCount, staleDeleted }
932}
933
934export type MediaCatalogueRecord = CatalogueRecord
935
936/**
937 * Full-sync media items (articles, podcasts, webinars, ...) into the
938 * media-catalogue index. Removed content is deleted from the store.
939 */
940export async function upsertMediaCatalogue(
941 records: MediaCatalogueRecord[],
942): Promise<{ totalVectors: number; staleDeleted: number }> {
943 return syncCatalogueIndex(mediaIndexName(), records, { idPrefix: 'media::' })
944}
945
946/**
947 * Full-sync site pages into the sitemap index so the agent can point users to
948 * the right URL. Removed pages are deleted from the store.
949 */
950export async function upsertSitemapCatalogue(
951 records: CatalogueRecord[],
952): Promise<{ totalVectors: number; staleDeleted: number }> {
953 return syncCatalogueIndex(sitemapIndexName(), records, { idPrefix: 'page::' })
954}
955
956/**
957 * Incremental sitemap sync for ONE page (publish webhook): embed + upsert the
958 * page's fresh records, then delete stale ids under `idPrefix`. Pass the
959 * page-scoped `<pageId>::` prefix — the trailing `::` keeps sibling pages
960 * safe (dashed slugs make one page id a plain prefix of a deeper page's id),
961 * and the meta vector (id == pageId, outside the prefix) is simply upserted.
962 */
963export async function syncSitemapPageVectors(
964 records: CatalogueRecord[],
965 idPrefix: string,
966): Promise<{ totalVectors: number; staleDeleted: number }> {
967 return syncCatalogueIndex(sitemapIndexName(), records, { idPrefix })
968}
969
970/**
971 * Full-sync Sanity-managed FAQs into the general store (SoW §8.1 / O03).
972 * FAQ vectors share the store with manually uploaded files — the faq:: id
973 * prefix keeps stale cleanup from ever touching the uploads.
974 */
975export async function upsertFaqCatalogue(
976 records: CatalogueRecord[],
977): Promise<{ totalVectors: number; staleDeleted: number }> {
978 return syncCatalogueIndex(scopeKeyToIndexName('general'), records, {
979 idPrefix: 'faq::',
980 })
981}
982
983/** Upsert or remove ONE catalogue record (incremental webhook updates). */
984export async function upsertCatalogueRecord(
985 indexName: string,
986 record: CatalogueRecord,
987): Promise<void> {
988 const openAiApiKey = requireEnv('OPENAI_API_KEY')
989 const [embedding] = await embedTexts(
990 [record.text],
991 openAiApiKey,
992 DEFAULTS.embeddingModel,
993 )
994 const index = getTargetIndex(indexName)
995 await upsertVectors(index, [
996 {
997 id: record.id,
998 values: embedding,
999 metadata: { ...record.metadata, text: record.text },
1000 },
1001 ])
1002 log('Catalogue record upserted', { indexName, id: record.id })
1003}
1004
1005export async function deleteCatalogueRecord(
1006 indexName: string,
1007 id: string,
1008): Promise<void> {
1009 const index = getTargetIndex(indexName)
1010 await withRetry('Pinecone delete record', () =>
1011 index.deleteMany({ ids: [id] }),
1012 )
1013 log('Catalogue record deleted', { indexName, id })
1014}
1015
1016export type StoreFile = { fileName: string; chunkCount: number }
1017
1018function assertValidScopeKey(scopeKey: string): void {
1019 if (
1020 !VALID_SCOPE_KEYS.includes(scopeKey as (typeof VALID_SCOPE_KEYS)[number])
1021 ) {
1022 throw new Error(
1023 `Invalid scopeKey. Must be one of: ${VALID_SCOPE_KEYS.join(', ')}`,
1024 )
1025 }
1026}
1027
1028/**
1029 * File inventory for a manual-upload store: every file name with its chunk
1030 * count, read from vector metadata (SoW §8.1: names + chunk counts visible in
1031 * Sanity). Upload stores are small, so listing + fetching is cheap.
1032 */
1033export async function listStoreFiles(scopeKey: string): Promise<StoreFile[]> {
1034 assertValidScopeKey(scopeKey)
1035 const index = getTargetIndex(scopeKeyToIndexName(scopeKey))
1036 const ids = [...(await listAllVectorIds(index))]
1037 const files = new Map<string, number>()
1038 for (const batch of sliceIntoBatches(ids, 100)) {
1039 const response = await withRetry('Pinecone fetch', () =>
1040 index.fetch({ ids: batch }),
1041 )
1042 for (const record of Object.values(response.records ?? {})) {
1043 const fileName = record.metadata?.fileName
1044 if (typeof fileName === 'string' && fileName) {
1045 files.set(fileName, (files.get(fileName) ?? 0) + 1)
1046 }
1047 }
1048 }
1049 return [...files.entries()]
1050 .map(([fileName, chunkCount]) => ({ fileName, chunkCount }))
1051 .sort((a, b) => a.fileName.localeCompare(b.fileName))
1052}
1053
1054/** Complete page bodies from the existing sitemap namespace, scoped by event IDs. */
1055export async function loadEventPageDocuments(event: 'paris' | 'miami') {
1056 const index = getTargetIndex(sitemapIndexName())
1057 const ids = [
1058 ...(await listAllVectorIds(index, `page::page::events-unleash-${event}`)),
1059 ...(await listAllVectorIds(index, `page::event::events-unleash-${event}`)),
1060 ]
1061 const pages = new Map<
1062 string,
1063 {
1064 slug: string
1065 title: string
1066 content: string
1067 vectorIds: string[]
1068 chunks: Array<{ position: number; text: string }>
1069 }
1070 >()
1071 for (const batch of sliceIntoBatches(ids, 100)) {
1072 const response = await withRetry('Pinecone fetch', () =>
1073 index.fetch({ ids: batch }),
1074 )
1075 for (const [id, record] of Object.entries(response.records ?? {})) {
1076 const slug = String(record.metadata?.slug ?? '')
1077 if (
1078 slug !== `/events/unleash-${event}` &&
1079 !slug.startsWith(`/events/unleash-${event}/`)
1080 )
1081 continue
1082 const page = pages.get(slug) ?? {
1083 slug,
1084 title: String(record.metadata?.title ?? ''),
1085 content: '',
1086 vectorIds: [],
1087 chunks: [],
1088 }
1089 page.vectorIds.push(id)
1090 const text = String(record.metadata?.text ?? '')
1091 const content = text.split('\nContent:\n')[1]
1092 if (content)
1093 page.chunks.push({
1094 position: Number(record.metadata?.chunk ?? 0),
1095 text: content,
1096 })
1097 pages.set(slug, page)
1098 }
1099 }
1100 return [...pages.values()].map(({ chunks, ...page }) => ({
1101 ...page,
1102 content: chunks
1103 .sort((a, b) => a.position - b.position)
1104 .map((chunk) => chunk.text)
1105 .join('\n'),
1106 }))
1107}
1108
1109/** Complete approved FAQ/upload text for exact policy lookups; no similarity ranking. */
1110export async function loadKnowledgeDocuments(
1111 scopeKey: string,
1112): Promise<Array<{ source: string; text: string }>> {
1113 assertValidScopeKey(scopeKey)
1114 const index = getTargetIndex(scopeKeyToIndexName(scopeKey))
1115 const ids = [...(await listAllVectorIds(index))]
1116 const files = new Map<string, Array<{ position: number; text: string }>>()
1117 for (const batch of sliceIntoBatches(ids, 100)) {
1118 const response = await withRetry('Pinecone fetch', () =>
1119 index.fetch({ ids: batch }),
1120 )
1121 for (const [id, record] of Object.entries(response.records ?? {})) {
1122 const text = String(record.metadata?.text ?? '')
1123 if (!text) continue
1124 const source = String(record.metadata?.fileName ?? id)
1125 const chunks = files.get(source) ?? []
1126 chunks.push({ position: Number(record.metadata?.chunkIndex ?? 0), text })
1127 files.set(source, chunks)
1128 }
1129 }
1130 return [...files].map(([source, chunks]) => ({
1131 source,
1132 text: chunks
1133 .sort((a, b) => a.position - b.position)
1134 .reduce((text, chunk) => {
1135 if (!text) return chunk.text
1136 let overlap = Math.min(CHUNK_OVERLAP, text.length, chunk.text.length)
1137 while (overlap > 0 && !text.endsWith(chunk.text.slice(0, overlap)))
1138 overlap--
1139 return text + (overlap ? chunk.text.slice(overlap) : `\n${chunk.text}`)
1140 }, ''),
1141 }))
1142}
1143
1144/** Delete a single uploaded file (all its chunks) from a store. */
1145export async function deleteStoreFile(
1146 scopeKey: string,
1147 fileName: string,
1148): Promise<{ deleted: number }> {
1149 assertValidScopeKey(scopeKey)
1150 const index = getTargetIndex(scopeKeyToIndexName(scopeKey))
1151 const ids = await findFileVectorIds(index, scopeKey, fileName)
1152 if (ids.length > 0) {
1153 for (const batch of sliceIntoBatches(ids, 1000)) {
1154 await withRetry('Pinecone delete file', () =>
1155 index.deleteMany({ ids: batch }),
1156 )
1157 }
1158 }
1159 log('File deleted from store', { scopeKey, fileName, deleted: ids.length })
1160 return { deleted: ids.length }
1161}
1162
1163export type PineconeSearchResult = {
1164 score: number
1165 text: string
1166 metadata: Record<string, unknown>
1167}
1168
1169/**
1170 * Semantic search against one store: embed the query (same model the stores
1171 * were built with) and return the top matches' stored text + metadata.
1172 */
1173export async function queryPineconeStore(
1174 indexName: string,
1175 query: string,
1176 topK = 6,
1177): Promise<PineconeSearchResult[]> {
1178 return queryPineconeStores([indexName], query, topK)
1179}
1180
1181// ---------------------------------------------------------------- agenda lookup
1182
1183export type AgendaSession = {
1184 title: string
1185 start: string
1186 end: string
1187 day: string
1188 /** "Wednesday 21 October 2026" as written by the agenda builder; '' on older uploads. */
1189 dayLabel: string
1190 room: string
1191 format: string
1192 keynote: boolean
1193 track: string
1194 access: string
1195 decisionPillarId?: string
1196 decisionPillar?: string
1197 topics: string
1198 speakers: string
1199 url: string
1200}
1201
1202const AGENDA_CACHE_TTL_MS = 10 * 60_000
1203const agendaCache = new Map<
1204 string,
1205 { sessions: AgendaSession[]; agendaPage: string | null; loadedAt: number }
1206>()
1207
1208/** The `- Agenda page:` line the markdown builder writes under the event heading. */
1209function parseAgendaPage(markdown: string): string | null {
1210 const match = markdown.match(/^- Agenda page: (\S+)$/m)
1211 return match ? canonicalAgendaUrl(match[1]) : null
1212}
1213
1214function field(block: string, name: string): string {
1215 const match = block.match(new RegExp(`^- ${name}: (.+)$`, 'm'))
1216 return match ? match[1].trim() : ''
1217}
1218
1219/** Parse the `### title` session blocks written by the agenda markdown builder. */
1220export function parseAgendaSessions(markdown: string): AgendaSession[] {
1221 const sessions: AgendaSession[] = []
1222 for (const raw of markdown.split(/\n(?=### )/)) {
1223 if (!raw.startsWith('### ')) continue
1224 const [head, ...rest] = raw.split('\n')
1225 const block = rest.join('\n')
1226 const time = field(block, 'Time')
1227 if (!time) continue
1228 const [start = '', end = ''] = time.split('->').map((t) => t.trim())
1229 sessions.push({
1230 title: head.slice(4).trim(),
1231 start,
1232 end,
1233 day: start.slice(0, 10),
1234 dayLabel: field(block, 'Day'),
1235 room: field(block, 'Room'),
1236 format: field(block, 'Format'),
1237 keynote: field(block, 'Format').toLowerCase() === 'keynote',
1238 track: field(block, 'Track / sub-event'),
1239 access: field(block, 'Access'),
1240 decisionPillarId: field(block, 'Decision pillar ID'),
1241 decisionPillar: field(block, 'Decision pillar'),
1242 topics: field(block, 'Topics'),
1243 speakers: field(block, 'Speakers'),
1244 url: canonicalAgendaUrl(field(block, 'Session URL')),
1245 })
1246 }
1247 return sessions
1248}
1249
1250/**
1251 * Every session of an event's agenda, read straight from the agenda store
1252 * (list + fetch, no embedding) so structural questions — which stages exist,
1253 * what is on Stage 2, which sessions are keynotes — are answered from the
1254 * complete set instead of the top-K of a similarity search. Cached per warm
1255 * instance; refreshed by the agenda upload.
1256 */
1257export async function loadAgendaSessions(
1258 scopeKey: 'event-world' | 'event-america',
1259): Promise<AgendaSession[]> {
1260 return (await loadAgenda(scopeKey)).sessions
1261}
1262
1263/**
1264 * The event's agenda page URL as recorded in the agenda store, so the
1265 * structured lookup can point visitors at the page where sessions are
1266 * browsed instead of the event homepage. Null until the agenda is uploaded.
1267 */
1268export async function loadAgendaPage(
1269 scopeKey: 'event-world' | 'event-america',
1270): Promise<string | null> {
1271 return (await loadAgenda(scopeKey)).agendaPage
1272}
1273
1274async function loadAgenda(
1275 scopeKey: 'event-world' | 'event-america',
1276): Promise<{ sessions: AgendaSession[]; agendaPage: string | null }> {
1277 const cached = agendaCache.get(scopeKey)
1278 if (cached && Date.now() - cached.loadedAt < AGENDA_CACHE_TTL_MS) {
1279 return cached
1280 }
1281 const index = getTargetIndex(scopeKeyToAgendaIndexName(scopeKey))
1282 const ids: string[] = []
1283 let paginationToken: string | undefined
1284 do {
1285 const page = await withRetry('Pinecone list', () =>
1286 index.listPaginated({ limit: 100, paginationToken }),
1287 )
1288 for (const vector of page.vectors ?? []) {
1289 if (vector.id) ids.push(vector.id)
1290 }
1291 paginationToken = page.pagination?.next
1292 } while (paginationToken)
1293
1294 const seen = new Set<string>()
1295 const sessions: AgendaSession[] = []
1296 let agendaPage: string | null = null
1297 for (const batch of sliceIntoBatches(ids, 100)) {
1298 const fetched = await withRetry('Pinecone fetch', () =>
1299 index.fetch({ ids: batch }),
1300 )
1301 for (const record of Object.values(fetched.records ?? {})) {
1302 // Agenda chunks are written by upsertTextContent with a string `text`.
1303 const text = String(record.metadata?.text ?? '')
1304 if (!text) continue
1305 agendaPage ??= parseAgendaPage(text)
1306 for (const session of parseAgendaSessions(text)) {
1307 const key = `${session.title}::${session.start}::${session.room}`
1308 if (seen.has(key)) continue
1309 seen.add(key)
1310 sessions.push(session)
1311 }
1312 }
1313 }
1314 sessions.sort((a, b) => a.start.localeCompare(b.start))
1315 const loaded = { sessions, agendaPage, loadedAt: Date.now() }
1316 agendaCache.set(scopeKey, loaded)
1317 return loaded
1318}
1319
1320export type AgendaLookupFilter = {
1321 stage?: string
1322 day?: string
1323 keynotesOnly?: boolean
1324 /** Person lookup: case-insensitive match on the speakers line ("Bersin", "Josh Bersin"). */
1325 speaker?: string
1326 /** Session lookup: case-insensitive match on the session title. */
1327 title?: string
1328 /**
1329 * Interest filter for personalised agendas: comma-separated terms
1330 * ("AI, skills, talent acquisition"); a session matches when ANY term
1331 * appears in its topics, track, title or format.
1332 */
1333 topics?: string
1334}
1335
1336/** Diacritic- and case-insensitive containment test for names and titles. */
1337function fuzzyIncludes(haystack: string, needle: string): boolean {
1338 const fold = (value: string) =>
1339 value
1340 .normalize('NFD')
1341 .replace(/[\u0300-\u036f]/g, '')
1342 .toLowerCase()
1343 .replace(/\s+/g, ' ')
1344 .trim()
1345 return fold(haystack).includes(fold(needle))
1346}
1347
1348const AGENDA_LOOKUP_MAX_SESSIONS = 40
1349
1350/**
1351 * Structured view of the agenda for the concierge: the stage list (every room
1352 * named "Stage N - …", with counts), the event days, and the sessions matching
1353 * the filter. Stage matching is forgiving ("stage 2", "worktech", "skills").
1354 */
1355export function summarizeAgenda(
1356 sessions: AgendaSession[],
1357 filter: AgendaLookupFilter,
1358) {
1359 const roomCounts = new Map<string, number>()
1360 for (const session of sessions) {
1361 if (!session.room) continue
1362 roomCounts.set(session.room, (roomCounts.get(session.room) ?? 0) + 1)
1363 }
1364 const stages = [...roomCounts]
1365 .filter(([room]) => /^stage \d+/i.test(room))
1366 .sort(([a], [b]) => a.localeCompare(b, undefined, { numeric: true }))
1367 .map(([name, sessionCount]) => ({ name, sessionCount }))
1368 const otherLocations = [...roomCounts.keys()]
1369 .filter((room) => !/^stage \d+/i.test(room))
1370 .sort()
1371 const days = [...new Set(sessions.map((s) => s.day).filter(Boolean))].sort()
1372
1373 const wanted = filter.stage?.trim().toLowerCase()
1374 const wantedNumber = wanted?.match(/\d+/)?.[0]
1375 let matched = sessions
1376 if (wanted) {
1377 matched = matched.filter((s) => {
1378 const room = s.room.toLowerCase()
1379 if (wantedNumber && /^stage \d+/.test(room)) {
1380 return room.match(/^stage (\d+)/)?.[1] === wantedNumber
1381 }
1382 return room.includes(wanted)
1383 })
1384 }
1385 if (filter.day) matched = matched.filter((s) => s.day === filter.day)
1386 if (filter.keynotesOnly)
1387 matched = matched.filter((s) => s.format.trim().toLowerCase() === 'keynote')
1388 const speaker = filter.speaker?.trim()
1389 if (speaker) {
1390 matched = matched.filter((s) => fuzzyIncludes(s.speakers, speaker))
1391 }
1392 const title = filter.title?.trim()
1393 if (title) matched = matched.filter((s) => fuzzyIncludes(s.title, title))
1394 const terms = (filter.topics ?? '')
1395 .split(/[,;/]|\band\b/i)
1396 .map((t) => t.trim())
1397 .filter((t) => t.length >= 2)
1398 if (terms.length > 0) {
1399 matched = matched.filter((s) => {
1400 const haystack = `${s.decisionPillar || ''} ${s.decisionPillarId || ''} ${s.topics} ${s.track} ${s.title} ${s.format}`
1401 return terms.some((term) => fuzzyIncludes(haystack, term))
1402 })
1403 }
1404
1405 const total = matched.length
1406 // An empty complete-session lookup proves absence from this agenda snapshot,
1407 // not absence from every speaker announcement or future programme update.
1408 const speakerLookup = speaker
1409 ? sessions.some((s) => fuzzyIncludes(s.speakers, speaker))
1410 ? { speaker, listed: true as const }
1411 : {
1412 speaker,
1413 listed: false as const,
1414 note:
1415 'No session in the current published agenda snapshot lists this person. This does not establish their announcement status. Say you cannot confirm they are speaking from this agenda and link the official speakers page for updates; do not claim they are definitely not announced.',
1416 }
1417 : undefined
1418 return {
1419 stages,
1420 otherLocations,
1421 days,
1422 totalSessions: sessions.length,
1423 matchedSessions: total,
1424 speakerLookup,
1425 truncated: total > AGENDA_LOOKUP_MAX_SESSIONS,
1426 sessions: matched.slice(0, AGENDA_LOOKUP_MAX_SESSIONS).map((s) => ({
1427 title: s.title,
1428 day: s.dayLabel || s.day,
1429 start: s.start,
1430 end: s.end,
1431 room: s.room,
1432 format: s.format,
1433 keynote: s.keynote,
1434 access: s.access,
1435 speakers: s.speakers,
1436 decisionPillarId: s.decisionPillarId,
1437 decisionPillar: s.decisionPillar,
1438 topics: s.topics,
1439 url: s.url || undefined,
1440 })),
1441 }
1442}
1443
1444/**
1445 * Semantic search fanned out over several stores with one embedding call.
1446 * Results are merged by score so a question asked on an event page can be
1447 * answered from the event store or the general FAQ store, whichever holds
1448 * it — the model no longer has to guess the right store first.
1449 */
1450export type StoreQuery = {
1451 indexName: string
1452 /** Pinecone metadata filter applied to this store only (e.g. scope a shared index to one event's pages). */
1453 filter?: Record<string, unknown>
1454}
1455
1456function normalizeStoreQuery(store: string | StoreQuery): StoreQuery {
1457 return typeof store === 'string' ? { indexName: store } : store
1458}
1459
1460export async function queryPineconeStores(
1461 indexNames: readonly (string | StoreQuery)[],
1462 query: string,
1463 topK = 6,
1464): Promise<PineconeSearchResult[]> {
1465 const seen = new Set<string>()
1466 const stores: StoreQuery[] = []
1467 for (const raw of indexNames) {
1468 const store = normalizeStoreQuery(raw)
1469 const key = `${store.indexName}::${JSON.stringify(store.filter ?? null)}`
1470 if (seen.has(key)) continue
1471 seen.add(key)
1472 stores.push(store)
1473 }
1474 if (stores.length === 0) return []
1475 const openAiApiKey = requireEnv('OPENAI_API_KEY')
1476
1477 const [embedding] = await embedTexts(
1478 [query],
1479 openAiApiKey,
1480 DEFAULTS.embeddingModel,
1481 )
1482
1483 // One store failing (DNS blip, index cold start) must not blank the whole
1484 // fan-out: keep what the other stores returned, throw only if none answered.
1485 const settled = await Promise.allSettled(
1486 stores.map(({ indexName, filter }) =>
1487 withRetry('Pinecone search', () =>
1488 getTargetIndex(indexName).query({
1489 vector: embedding,
1490 topK,
1491 includeValues: false,
1492 includeMetadata: true,
1493 ...(filter ? { filter } : {}),
1494 }),
1495 ),
1496 ),
1497 )
1498 const failures = settled.filter(
1499 (result): result is PromiseRejectedResult => result.status === 'rejected',
1500 )
1501 if (failures.length === settled.length) throw failures[0].reason
1502 for (const [i, result] of settled.entries()) {
1503 if (result.status === 'rejected') {
1504 console.warn(
1505 `${LOG_PREFIX} search on ${stores[i].indexName} failed; continuing with the other stores`,
1506 result.reason,
1507 )
1508 }
1509 }
1510 const perStore = settled.flatMap((result) =>
1511 result.status === 'fulfilled' ? [result.value] : [],
1512 )
1513
1514 return perStore
1515 .flatMap((response) => response.matches ?? [])
1516 .map((match) => {
1517 // Every store is written by this module with a string `text` field and
1518 // flat scalar metadata (RecordMetadata), so the split below is safe.
1519 const { text, ...rest } = match.metadata ?? {}
1520 return {
1521 score: match.score ?? 0,
1522 text: String(text ?? ''),
1523 metadata: rest,
1524 }
1525 })
1526 .sort((a, b) => b.score - a.score)
1527 .slice(0, topK)
1528}
1529