apps/web/src/routes/api/pinecone/sync.ts

1import {
2	mediaIndexName,
3	sitemapIndexName,
4	syncSitemapPageVectors,
5	upsertCatalogueRecord,
6	upsertFaqCatalogue,
7} from '@local/ai/components/api/pineconeUpload'
8import { EDITORIAL_TYPES } from '@local/config'
9import { getPublishedSanityClient } from '@local/sanity/client'
10import { sanityFetch } from '@local/sanity/utils/fetch'
11import { isValidSignature, SIGNATURE_HEADER_NAME } from '@sanity/webhook'
12import type { APIEvent } from '@solidjs/start/server'
13import { faqCatalogueRecords, fetchFaqsFromSanity } from '~/lib/faqCatalogue'
14import { formatMediaItemText, mediaItemId } from '~/lib/media'
15import {
16	enrichSponsorReferences,
17	eventPageContentRecords,
18	eventPageMetaRecord,
19	extractPageOutline,
20	extractSlicesText,
21	formatSitePageText,
22	isEventSitePage,
23	loadSponsorReferenceTitles,
24	pageTitle,
25	pageVectorMetadata,
26	sitePageId,
27} from '~/lib/pageCatalogue'
28
29/**
30 * Sanity publish webhook → vector store freshness (SoW §8.1 / O03).
31 * Register in sanity.io/manage with the REVALIDATE_SECRET (see docs/ai-agent-ops.md):
32 *  - faqSet changes trigger a full FAQ resync (cheap; always correct)
33 *  - media/page documents get an incremental single-vector upsert, or a
34 *    delete when the document is gone
35 * Nightly full refreshes remain the safety net for anything incremental misses.
36 */
37
38const LOG_PREFIX = '[Pinecone sync]'
39
40function log(msg: string, data?: Record<string, unknown>) {
41	if (process.env.SYNC_DEBUG !== '1' && process.env.SYNC_DEBUG !== 'true') return
42	const payload = data ? ` ${JSON.stringify(data)}` : ''
43	console.log(`${LOG_PREFIX} ${msg}${payload}`)
44}
45
46const MEDIA_TYPES = new Set<string>([...EDITORIAL_TYPES, 'podcast', 'webinar'])
47const PAGE_TYPES = new Set<string>(['page', 'event', 'taxonomy-page', 'home'])
48
49type WebhookPayload = {
50	_id?: string
51	_type?: string
52}
53
54async function syncFaqs() {
55	const entries = await fetchFaqsFromSanity()
56	const records = faqCatalogueRecords(entries)
57	const { totalVectors, staleDeleted } = await upsertFaqCatalogue(records)
58	return { action: 'faq-resync', totalVectors, staleDeleted }
59}
60
61/** Fetch the published doc; upsert its vector, or delete it if unpublished. */
62async function syncMediaDoc(id: string) {
63	const doc = await sanityFetch<{
64		_type: string
65		slug?: string
66		title?: string
67		excerpt?: string
68		topics?: string[]
69		author?: string
70	} | null>(
71		getPublishedSanityClient(),
72		`*[_id == $id][0]{_type, "slug": coalesce(fullSlug, slug.fullUrl, slug.current, ""), title, excerpt, "topics": topics[]->title, "author": author->title}`,
73		{ id },
74		{ tag: 'page-catalogue' },
75	)
76	if (!doc || !doc.slug) {
77		// Unpublished/deleted and we can't compute the vector id without the
78		// slug: rely on the next full refresh. Deletions with a known doc still
79		// present as drafts resolve on republish.
80		log('media doc gone; awaiting full refresh for cleanup', { id })
81		return { action: 'media-skip', id }
82	}
83	const item = {
84		type: doc._type,
85		slug: doc.slug,
86		title: doc.title?.trim() || 'Untitled',
87		excerpt: doc.excerpt?.trim() || null,
88		topics: (doc.topics ?? []).filter((t): t is string => typeof t === 'string'),
89		author: typeof doc.author === 'string' ? doc.author : null,
90		metaLine: null,
91		id,
92	}
93	await upsertCatalogueRecord(mediaIndexName(), {
94		id: mediaItemId(item.type, item.slug),
95		text: formatMediaItemText(item),
96		metadata: { type: item.type, slug: item.slug, title: item.title },
97	})
98	return { action: 'media-upsert', slug: doc.slug }
99}
100
101async function syncPageDoc(id: string) {
102	// Slices are only fetched for event pages — the only ones whose body copy
103	// is indexed as content chunks; every other page stays meta-only.
104	const doc = await sanityFetch<{
105		_type: string
106		slug?: string
107		agentTitle?: string
108		title?: string
109		description?: string
110		slices?: unknown
111	} | null>(
112		getPublishedSanityClient(),
113		`*[_id == $id][0]{_type, "slug": fullSlug, agentTitle, title, "description": coalesce(agentDescription, seo.metaDescription, seo.description, excerpt), "slices": select(_type in ["page", "event"] && string::startsWith(fullSlug, "/events/") => slices, null)}`,
114		{ id },
115		{ tag: 'page-catalogue' },
116	)
117	if (!doc || !doc.slug) {
118		log('page doc gone; awaiting full refresh for cleanup', { id })
119		return { action: 'page-skip', id }
120	}
121	const page = {
122		type: doc._type,
123		slug: doc.slug,
124		title: pageTitle(doc, doc.slug),
125		description: doc.description?.trim() || null,
126	}
127	const metaRecord = {
128		id: sitePageId(page),
129		text: formatSitePageText(page),
130		metadata: pageVectorMetadata(page),
131	}
132	if (isEventSitePage(page)) {
133		const sponsorTitles = await loadSponsorReferenceTitles(doc.slices)
134		const eventPage = {
135			...page,
136			content: extractSlicesText(
137				enrichSponsorReferences(doc.slices, sponsorTitles),
138			),
139			outline: extractPageOutline(doc.slices),
140		}
141		const records = [
142			eventPageMetaRecord(eventPage),
143			...eventPageContentRecords(eventPage),
144		]
145		// Page-scoped `<pageId>::` prefix: shrinking content drops stale tail
146		// chunks without ever touching other pages' vectors.
147		const { totalVectors, staleDeleted } = await syncSitemapPageVectors(
148			records,
149			`${sitePageId(page)}::`,
150		)
151		return {
152			action: 'page-content-upsert',
153			slug: doc.slug,
154			totalVectors,
155			staleDeleted,
156		}
157	}
158	await upsertCatalogueRecord(sitemapIndexName(), metaRecord)
159	return { action: 'page-upsert', slug: doc.slug }
160}
161
162export async function POST({ request }: APIEvent) {
163	const signature = request.headers.get(SIGNATURE_HEADER_NAME)
164	const rawBody = await request.text()
165	const isValid = await isValidSignature(
166		rawBody,
167		signature ?? '',
168		process.env.REVALIDATE_SECRET ?? '',
169	)
170	if (!isValid) {
171		return Response.json(
172			{ ok: false, error: 'Invalid signature' },
173			{ status: 401 },
174		)
175	}
176
177	let payload: WebhookPayload
178	try {
179		payload = JSON.parse(rawBody) as WebhookPayload
180	} catch {
181		return Response.json({ ok: false, error: 'Invalid JSON' }, { status: 400 })
182	}
183
184	const type = payload._type ?? ''
185	const id = payload._id?.replace(/^drafts\./, '') ?? ''
186	log('Webhook received', { type, id })
187
188	try {
189		let result: Record<string, unknown> = { action: 'ignored', type }
190		if (type === 'faqSet') {
191			result = await syncFaqs()
192		} else if (MEDIA_TYPES.has(type) && id) {
193			result = await syncMediaDoc(id)
194			// Media docs are pages too — keep the sitemap entry fresh as well.
195			await syncPageDoc(id).catch(() => undefined)
196		} else if (PAGE_TYPES.has(type) && id) {
197			result = await syncPageDoc(id)
198		} else if ((type === 'company' || type === 'profile') && id) {
199			const pages = await sanityFetch<Array<{ _id: string }>>(
200				getPublishedSanityClient(),
201				'*[_type in ["page", "event"] && references($id)]{_id}',
202				{ id },
203				{ tag: 'page-catalogue' },
204			)
205			for (const page of pages) await syncPageDoc(page._id)
206			result = { action: 'referenced-page-resync', pages: pages.length }
207		}
208		log('Sync done', result)
209		return Response.json({ ok: true, ...result })
210	} catch (error) {
211		const message = error instanceof Error ? error.message : 'Sync failed'
212		console.error(`${LOG_PREFIX} failed`, error)
213		return Response.json({ ok: false, error: message }, { status: 500 })
214	}
215}
216