2026-04-12 22:21:39 -04:00
|
|
|
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<string> {
|
|
|
|
|
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, any>): 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;
|
|
|
|
|
}
|
|
|
|
|
|
2026-06-08 15:42:40 -04:00
|
|
|
// 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);
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-12 22:21:39 -04:00
|
|
|
async function clearOldFiles(bucket: string): Promise<void> {
|
|
|
|
|
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<Record<string, any>[]> {
|
|
|
|
|
const results: Record<string, any>[] = [];
|
|
|
|
|
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<string, any>[]));
|
|
|
|
|
cursor = res.has_more && res.next_cursor ? res.next_cursor : undefined;
|
|
|
|
|
} while (cursor);
|
|
|
|
|
|
|
|
|
|
return results;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
export const handler = async (): Promise<void> => {
|
|
|
|
|
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,
|
2026-06-08 15:42:40 -04:00
|
|
|
'notion-title': sanitizeMetadataValue(title),
|
2026-04-12 22:21:39 -04:00
|
|
|
},
|
|
|
|
|
}),
|
|
|
|
|
);
|
|
|
|
|
|
|
|
|
|
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...');
|
2026-06-08 15:42:40 -04:00
|
|
|
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;
|
|
|
|
|
}
|
2026-04-12 22:21:39 -04:00
|
|
|
};
|