import type { RealtimeChannel } from '@supabase/supabase-js' import type { D3roSupabaseClient, KnowledgeDocument, KnowledgeFileType, } from '@d3ro/api-client' import { SUPABASE_URL } from '@d3ro/core/supabase-config' import { supabase } from '../../lib/supabase' import { createUuidV4 } from '../../lib/random-id' import { type ForegroundReconciliationDeps, knowledgeRealtimeBindings, startForegroundReconciliation, } from './knowledge-realtime' export type MobileKnowledgeFileType = Extract export type KnowledgeServiceErrorCode = | 'auth' | 'cancelled' | 'conflict' | 'forbidden' | 'index-failed' | 'index-unavailable' | 'invalid-response' | 'network' | 'not-found' | 'server' | 'timeout' | 'validation' export class KnowledgeServiceError extends Error { constructor( public readonly code: KnowledgeServiceErrorCode, public readonly retryable: boolean, public readonly status: number | null = null, message: string = code, ) { super(message) this.name = 'KnowledgeServiceError' } } export interface KnowledgeListOptions { userId: string search?: string page?: number pageSize?: number } export interface KnowledgeListResult { documents: KnowledgeDocument[] total: number page: number pageSize: number hasMore: boolean } export interface CreateKnowledgeDocumentInput { userId: string title: string fileName: string fileType: MobileKnowledgeFileType content: string } export interface KnowledgeSearchResult { id: string documentId: string chunkIndex: number content: string similarity: number } export interface KnowledgeDocumentSubscription { unsubscribe: () => Promise } interface KnowledgeSearchOptions { accessToken: string query: string count?: number signal?: AbortSignal timeoutMs?: number } interface KnowledgeIndexOptions { accessToken: string userId: string document: KnowledgeDocument signal?: AbortSignal timeoutMs?: number } const UUID_PATTERN = /^[0-9a-f]{8}-[0-9a-f]{4}-[1-5][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i const DEFAULT_PAGE_SIZE = 30 const MAX_PAGE_SIZE = 50 const MAX_TITLE_LENGTH = 160 export const KNOWLEDGE_MAX_CONTENT_CHARS = 250_000 export const KNOWLEDGE_CHUNK_SIZE = 800 const MAX_SEARCH_LENGTH = 500 const MAX_SEARCH_RESULTS = 12 function client(): D3roSupabaseClient { return supabase as unknown as D3roSupabaseClient } function isRecord(value: unknown): value is Record { return typeof value === 'object' && value !== null && !Array.isArray(value) } function requireUuid(value: string, label: string): void { if (!UUID_PATTERN.test(value)) { throw new KnowledgeServiceError('auth', false, null, `A valid ${label} is required`) } } export function normalizeKnowledgeTitle(value: string): string { const title = value.trim().replace(/\s+/g, ' ') if (title.length === 0 || title.length > MAX_TITLE_LENGTH) { throw new KnowledgeServiceError( 'validation', false, null, `Knowledge titles must contain 1-${MAX_TITLE_LENGTH} characters`, ) } return title } export function sanitizeKnowledgeSearch(value: string): string { return value .trim() .slice(0, MAX_SEARCH_LENGTH) .replace(/[\\%_,()]/g, ' ') .replace(/\s+/g, ' ') .trim() } export function normalizeKnowledgeContent(value: string): string { const withoutBom = value.charCodeAt(0) === 0xfeff ? value.slice(1) : value const content = withoutBom.trim() if ( content.length === 0 || content.length > KNOWLEDGE_MAX_CONTENT_CHARS || content.includes('\u0000') ) { throw new KnowledgeServiceError( 'validation', false, null, `Knowledge content must contain 1-${KNOWLEDGE_MAX_CONTENT_CHARS} text characters`, ) } return content } export function chunkKnowledgeText(value: string): string[] { const content = normalizeKnowledgeContent(value) const chunks: string[] = [] let offset = 0 while (offset < content.length) { const hardEnd = Math.min(offset + KNOWLEDGE_CHUNK_SIZE, content.length) let end = hardEnd if (hardEnd < content.length) { const candidate = content.lastIndexOf('\n', hardEnd) if (candidate > offset + Math.floor(KNOWLEDGE_CHUNK_SIZE * 0.6)) end = candidate } const chunk = content.slice(offset, end).trim() if (chunk.length > 0) chunks.push(chunk) offset = end while (offset < content.length && /\s/.test(content[offset])) offset += 1 } if (chunks.length === 0) { throw new KnowledgeServiceError('validation', false, null, 'Knowledge content is empty') } return chunks } function assertDocument(row: KnowledgeDocument | null, expectedId?: string): KnowledgeDocument { if (row === null) { throw new KnowledgeServiceError('not-found', false, 404, 'Knowledge document was not found') } requireUuid(row.id, 'knowledge document id') requireUuid(row.user_id, 'knowledge owner id') if (expectedId !== undefined && row.id !== expectedId) { throw new KnowledgeServiceError('invalid-response', true, null, 'Server returned a different document') } if ( typeof row.title !== 'string' || row.title.trim().length === 0 || row.title.length > 500 || (row.file_name !== null && (typeof row.file_name !== 'string' || row.file_name.length > 255)) || typeof row.indexed !== 'boolean' || !Number.isSafeInteger(row.chunk_count) || row.chunk_count < 0 || !Number.isFinite(Date.parse(row.created_at)) || !Number.isFinite(Date.parse(row.updated_at)) ) { throw new KnowledgeServiceError('invalid-response', true, null, 'Server returned an invalid document') } return row } function toKnowledgeError(error: unknown): KnowledgeServiceError { if (error instanceof KnowledgeServiceError) return error const candidate = error as { code?: unknown; message?: unknown } const code = typeof candidate?.code === 'string' ? candidate.code : '' const message = typeof candidate?.message === 'string' ? candidate.message : 'Knowledge request failed' const lower = message.toLowerCase() if ( error instanceof TypeError || lower.includes('network request failed') || lower.includes('failed to fetch') || lower.includes('networkerror') ) { return new KnowledgeServiceError('network', true, null, message) } if (code === '42501' || lower.includes('permission denied') || lower.includes('row-level security')) { return new KnowledgeServiceError('forbidden', false, 403, message) } if (code === 'PGRST301' || lower.includes('jwt')) { return new KnowledgeServiceError('auth', false, 401, message) } if (code.startsWith('22')) { return new KnowledgeServiceError('validation', false, 400, message) } return new KnowledgeServiceError('server', true, null, message) } function createRequestDeadline( externalSignal: AbortSignal | undefined, timeoutMs: number, ): { signal: AbortSignal; dispose: () => void; didTimeout: () => boolean } { const controller = new AbortController() let timedOut = false const abortFromCaller = (): void => controller.abort() if (externalSignal?.aborted === true) controller.abort() else externalSignal?.addEventListener('abort', abortFromCaller, { once: true }) const timer = setTimeout(() => { timedOut = true controller.abort() }, timeoutMs) return { signal: controller.signal, didTimeout: () => timedOut, dispose: () => { clearTimeout(timer) externalSignal?.removeEventListener('abort', abortFromCaller) }, } } async function responseJson(response: Response): Promise { return response.json().catch(() => null) } function edgeErrorCode(status: number, body: unknown): KnowledgeServiceErrorCode { const serverCode = isRecord(body) && typeof body.error === 'string' ? body.error : '' if (status === 401) return 'auth' if (status === 400) return 'validation' if (status === 403) return 'forbidden' if (status === 404) return 'not-found' if ( status === 503 && (serverCode === 'embedding_provider_unavailable' || serverCode === 'knowledge_storage_unavailable') ) return 'index-unavailable' if ( status === 502 && (serverCode === 'embedding_failed' || serverCode === 'embedding_upstream_failed' || serverCode === 'embedding_response_invalid') ) return 'index-failed' if (status === 504) return 'timeout' return 'server' } export async function listKnowledgeDocuments( options: KnowledgeListOptions, ): Promise { requireUuid(options.userId, 'authenticated user') const page = Math.max(0, Math.floor(options.page ?? 0)) const pageSize = Math.max(1, Math.min(Math.floor(options.pageSize ?? DEFAULT_PAGE_SIZE), MAX_PAGE_SIZE)) const search = sanitizeKnowledgeSearch(options.search ?? '') try { let query = client() .from('knowledge_documents') .select('*', { count: 'exact' }) .order('created_at', { ascending: false }) .order('id', { ascending: false }) .range(page * pageSize, (page + 1) * pageSize - 1) if (search.length > 0) query = query.ilike('title', `%${search}%`) const { data, error, count } = await query if (error !== null) throw error const documents = (data ?? []).map((row) => assertDocument(row)) const total = count ?? documents.length return { documents, total, page, pageSize, hasMore: (page + 1) * pageSize < total, } } catch (error) { throw toKnowledgeError(error) } } export function mergeKnowledgeDocuments( current: KnowledgeDocument[], incoming: KnowledgeDocument[], ): KnowledgeDocument[] { const merged = new Map(current.map((document) => [document.id, document])) for (const document of incoming) { const existing = merged.get(document.id) if (existing === undefined || document.updated_at >= existing.updated_at) { merged.set(document.id, document) } } return [...merged.values()].sort((left, right) => { const byCreated = right.created_at.localeCompare(left.created_at) return byCreated !== 0 ? byCreated : right.id.localeCompare(left.id) }) } export async function createKnowledgeDocument( input: CreateKnowledgeDocumentInput, ): Promise { requireUuid(input.userId, 'authenticated user') const title = normalizeKnowledgeTitle(input.title) const fileName = input.fileName.trim() if (fileName.length === 0 || fileName.length > 255 || !['txt', 'md'].includes(input.fileType)) { throw new KnowledgeServiceError('validation', false, null, 'Knowledge file metadata is invalid') } const chunks = chunkKnowledgeText(input.content) let created: KnowledgeDocument | null = null try { const documentResult = await client() .from('knowledge_documents') .insert({ user_id: input.userId, title, file_name: fileName, file_type: input.fileType, storage_key: null, chunk_count: chunks.length, indexed: false, indexed_at: null, }) .select('*') .single() if (documentResult.error !== null) throw documentResult.error created = assertDocument(documentResult.data) const chunkResult = await client() .from('knowledge_chunks') .insert(chunks.map((content, chunkIndex) => ({ document_id: created?.id ?? '', chunk_index: chunkIndex, content, }))) .select('id') if (chunkResult.error !== null) throw chunkResult.error if ((chunkResult.data ?? []).length !== chunks.length) { throw new KnowledgeServiceError( 'invalid-response', true, null, 'Server did not confirm every knowledge chunk', ) } return created } catch (error) { if (created !== null) { const cleanup = await client() .from('knowledge_documents') .delete() .eq('id', created.id) .eq('user_id', input.userId) if (cleanup.error !== null) { throw new KnowledgeServiceError( 'server', true, null, 'Knowledge import failed and its incomplete document could not be removed', ) } } throw toKnowledgeError(error) } } async function forcePendingIndexState( userId: string, documentId: string, ): Promise { await client() .from('knowledge_documents') .update({ indexed: false, indexed_at: null }) .eq('id', documentId) .eq('user_id', userId) } export async function indexKnowledgeDocument( options: KnowledgeIndexOptions, ): Promise { requireUuid(options.userId, 'authenticated user') const document = assertDocument(options.document) if (document.user_id !== options.userId) { throw new KnowledgeServiceError('forbidden', false, 403, 'Only the document owner can index it') } const token = options.accessToken.trim() if (token.length === 0) throw new KnowledgeServiceError('auth', false, 401) const deadline = createRequestDeadline(options.signal, options.timeoutMs ?? 90_000) try { const response = await fetch(`${SUPABASE_URL}/functions/v1/embed-chunks`, { method: 'POST', headers: { Authorization: `Bearer ${token}`, 'Content-Type': 'application/json', }, body: JSON.stringify({ document_id: document.id }), signal: deadline.signal, }) const body = await responseJson(response) if (!response.ok) { const code = edgeErrorCode(response.status, body) throw new KnowledgeServiceError( code, code === 'index-unavailable' || code === 'timeout' || code === 'server', response.status, ) } if ( !isRecord(body) || !Number.isSafeInteger(body.embedded) || Number(body.embedded) < 0 || !Number.isSafeInteger(body.total) || Number(body.total) !== document.chunk_count || body.indexed !== true ) { throw new KnowledgeServiceError( 'invalid-response', true, response.status, 'Embedding service did not confirm a complete index', ) } const [documentResult, embeddedCountResult] = await Promise.all([ client() .from('knowledge_documents') .select('*') .eq('id', document.id) .maybeSingle(), client() .from('knowledge_chunks') .select('id', { count: 'exact', head: true }) .eq('document_id', document.id) .not('embedding', 'is', null), ]) if (documentResult.error !== null) throw documentResult.error if (embeddedCountResult.error !== null) throw embeddedCountResult.error const updated = assertDocument(documentResult.data, document.id) if (updated.indexed !== true || embeddedCountResult.count !== updated.chunk_count) { throw new KnowledgeServiceError( 'invalid-response', true, response.status, 'Server marked an incomplete knowledge index as ready', ) } return updated } catch (error) { await forcePendingIndexState(options.userId, document.id).catch(() => undefined) if (error instanceof KnowledgeServiceError) throw error if (deadline.signal.aborted) { const callerCancelled = options.signal?.aborted === true && !deadline.didTimeout() throw new KnowledgeServiceError( callerCancelled ? 'cancelled' : 'timeout', !callerCancelled, callerCancelled ? null : 504, ) } throw toKnowledgeError(error) } finally { deadline.dispose() } } export async function deleteKnowledgeDocument( userId: string, document: KnowledgeDocument, ): Promise { requireUuid(userId, 'authenticated user') assertDocument(document) if (document.user_id !== userId) { throw new KnowledgeServiceError('forbidden', false, 403, 'Only the owner can delete this document') } try { const { data, error } = await client() .from('knowledge_documents') .delete() .eq('id', document.id) .eq('user_id', userId) .eq('updated_at', document.updated_at) .select('id') .maybeSingle() if (error !== null) throw error if (data === null) { throw new KnowledgeServiceError( 'conflict', true, 409, 'Knowledge document changed or was removed on another device', ) } } catch (error) { throw toKnowledgeError(error) } } function parseSearchResult(value: unknown): KnowledgeSearchResult { if (!isRecord(value)) throw new KnowledgeServiceError('invalid-response', true) const id = typeof value.id === 'string' ? value.id : '' const documentId = typeof value.document_id === 'string' ? value.document_id : '' const content = typeof value.content === 'string' ? value.content.trim() : '' const chunkIndex = value.chunk_index const similarity = value.similarity if ( !UUID_PATTERN.test(id) || !UUID_PATTERN.test(documentId) || content.length === 0 || content.length > 8_000 || !Number.isSafeInteger(chunkIndex) || Number(chunkIndex) < 0 || typeof similarity !== 'number' || !Number.isFinite(similarity) || similarity < 0 || similarity > 1 ) { throw new KnowledgeServiceError('invalid-response', true) } return { id, documentId, chunkIndex: Number(chunkIndex), content, similarity, } } export async function searchKnowledge( options: KnowledgeSearchOptions, ): Promise { const token = options.accessToken.trim() if (token.length === 0) throw new KnowledgeServiceError('auth', false, 401) const query = options.query.trim() if (query.length === 0 || query.length > MAX_SEARCH_LENGTH) { throw new KnowledgeServiceError('validation', false, 400) } const count = Math.max(1, Math.min(Math.floor(options.count ?? 8), MAX_SEARCH_RESULTS)) const deadline = createRequestDeadline(options.signal, options.timeoutMs ?? 45_000) try { const response = await fetch(`${SUPABASE_URL}/functions/v1/search-knowledge`, { method: 'POST', headers: { Authorization: `Bearer ${token}`, 'Content-Type': 'application/json', }, body: JSON.stringify({ query, count }), signal: deadline.signal, }) const body = await responseJson(response) if (!response.ok) { const code = edgeErrorCode(response.status, body) throw new KnowledgeServiceError( code, code === 'index-unavailable' || code === 'timeout' || code === 'server', response.status, ) } if (!isRecord(body) || !Array.isArray(body.results) || body.results.length > count) { throw new KnowledgeServiceError('invalid-response', true, response.status) } return body.results.map(parseSearchResult) } catch (error) { if (error instanceof KnowledgeServiceError) throw error if (deadline.signal.aborted) { const callerCancelled = options.signal?.aborted === true && !deadline.didTimeout() throw new KnowledgeServiceError( callerCancelled ? 'cancelled' : 'timeout', !callerCancelled, callerCancelled ? null : 504, ) } throw toKnowledgeError(error) } finally { deadline.dispose() } } export interface KnowledgeDocumentSubscriptionOptions { /** * Owner whose documents the list shows. When given, INSERT/UPDATE events are * narrowed to that owner so teammates' shared-document writes do not trigger * list refetches for a list that only shows owned documents. */ userId?: string /** Injectable foreground/interval ports for deletion reconciliation. */ reconciliation?: ForegroundReconciliationDeps } /** * Live updates for the Knowledge list. Only INSERT/UPDATE are subscribed: * Postgres Changes cannot filter DELETE events and does not apply RLS to them, * so a DELETE subscription would leak other accounts' document ids and fan out * a refetch to every client on every deletion. Deletions are reconciled by * foreground polling instead (same policy as use-history-sync.ts). */ export function subscribeToKnowledgeDocuments( onChanged: () => void, onStatus: (connected: boolean) => void, options: KnowledgeDocumentSubscriptionOptions = {}, ): KnowledgeDocumentSubscription { if (options.userId !== undefined) requireUuid(options.userId, 'knowledge owner id') let channel: RealtimeChannel = supabase.channel(`mobile-knowledge-${createUuidV4()}`) for (const binding of knowledgeRealtimeBindings(options.userId)) { channel = channel.on('postgres_changes', binding, onChanged) } channel = channel.subscribe((status) => onStatus(status === 'SUBSCRIBED')) const reconciliation = startForegroundReconciliation(onChanged, options.reconciliation) return { unsubscribe: async () => { reconciliation.stop() await supabase.removeChannel(channel) }, } }