import { Client } from '@notionhq/client'; import { NotionToMarkdown } from 'notion-to-md'; import { S3Client, PutObjectCommand, ListObjectsV2Command, DeleteObjectsCommand, } from '@aws-sdk/client-s3'; import { BedrockAgentClient, StartIngestionJobCommand, } from '@aws-sdk/client-bedrock-agent'; import { SecretsManagerClient, GetSecretValueCommand, } from '@aws-sdk/client-secrets-manager'; const s3 = new S3Client({}); const bedrockAgent = new BedrockAgentClient({ region: process.env.REGION! }); const secrets = new SecretsManagerClient({}); const NOTION_PREFIX = 'notion/'; let cachedApiKey: string | undefined; async function getApiKey(): Promise { if (cachedApiKey) return cachedApiKey; const res = await secrets.send( new GetSecretValueCommand({ SecretId: process.env.NOTION_SECRET_ARN! }), ); const parsed = JSON.parse(res.SecretString!); cachedApiKey = parsed.apiKey; return cachedApiKey!; } function getPageTitle(page: Record): string { // Regular pages store the title in properties.title or properties.Name const props = page.properties ?? {}; for (const prop of Object.values(props) as any[]) { if (prop?.type === 'title' && Array.isArray(prop.title) && prop.title[0]?.plain_text) { return prop.title[0].plain_text; } } return page.id as string; } // S3 user-metadata values travel as HTTP headers, which must be US-ASCII with no // control characters. Notion titles routinely contain em dashes and other // non-ASCII characters (ERR_INVALID_CHAR otherwise). Strip to safe printable // ASCII and cap at the 256-char metadata limit. function sanitizeMetadataValue(value: string): string { return value .replace(/[^\x20-\x7E]/g, '') // drop non-printable / non-ASCII .trim() .slice(0, 256); } async function clearOldFiles(bucket: string): Promise { let continuationToken: string | undefined; const toDelete: { Key: string }[] = []; do { const res = await s3.send( new ListObjectsV2Command({ Bucket: bucket, Prefix: NOTION_PREFIX, ContinuationToken: continuationToken, }), ); for (const obj of res.Contents ?? []) { if (obj.Key) toDelete.push({ Key: obj.Key }); } continuationToken = res.NextContinuationToken; } while (continuationToken); if (toDelete.length === 0) return; // S3 DeleteObjects accepts up to 1000 keys per request for (let i = 0; i < toDelete.length; i += 1000) { await s3.send( new DeleteObjectsCommand({ Bucket: bucket, Delete: { Objects: toDelete.slice(i, i + 1000) }, }), ); } console.log(`Deleted ${toDelete.length} stale Notion file(s) from S3`); } async function getAllPages(notion: Client): Promise[]> { const results: Record[] = []; let cursor: string | undefined; do { const res = await notion.search({ filter: { value: 'page', property: 'object' }, page_size: 100, start_cursor: cursor, }); results.push(...(res.results as Record[])); cursor = res.has_more && res.next_cursor ? res.next_cursor : undefined; } while (cursor); return results; } export const handler = async (): Promise => { const apiKey = await getApiKey(); const notion = new Client({ auth: apiKey }); const n2m = new NotionToMarkdown({ notionClient: notion }); const bucket = process.env.KB_BUCKET_NAME!; const kbId = process.env.KNOWLEDGE_BASE_ID!; const dsId = process.env.DATA_SOURCE_ID!; console.log('Clearing stale Notion files from S3...'); await clearOldFiles(bucket); console.log('Fetching pages from Notion Office Operations teamspace...'); const pages = await getAllPages(notion); console.log(`Found ${pages.length} page(s)`); let synced = 0; let failed = 0; for (const page of pages) { const pageId = page.id as string; const title = getPageTitle(page); try { const mdBlocks = await n2m.pageToMarkdown(pageId); const { parent: mdContent } = n2m.toMarkdownString(mdBlocks); await s3.send( new PutObjectCommand({ Bucket: bucket, Key: `${NOTION_PREFIX}${pageId}.md`, Body: `# ${title}\n\n${mdContent}`, ContentType: 'text/markdown', Metadata: { 'notion-page-id': pageId, 'notion-title': sanitizeMetadataValue(title), }, }), ); console.log(`Synced: "${title}" (${pageId})`); synced++; } catch (err) { console.error(`Failed to sync page "${title}" (${pageId}):`, err); failed++; } } console.log(`Upload complete. ${synced} synced, ${failed} failed.`); console.log('Triggering Bedrock KB ingestion job...'); try { const ingestionRes = await bedrockAgent.send( new StartIngestionJobCommand({ knowledgeBaseId: kbId, dataSourceId: dsId, }), ); console.log( `Ingestion job started: ${ingestionRes.ingestionJob?.ingestionJobId}`, ); } catch (err) { // A ConflictException means an ingestion job is already running for this data // source (e.g. an overlapping run). The freshly uploaded files will be picked // up by that in-flight job, so this is benign — log and exit cleanly rather // than failing the whole invocation. if (err instanceof Error && err.name === 'ConflictException') { console.log( 'Ingestion job already in progress; uploaded files will be picked up by the running job. Skipping.', ); return; } throw err; } };