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