/** * Copyright 2023 Kapeta Inc. * SPDX-License-Identifier: BUSL-1.1 */ import Router from 'express-promise-router'; import FS from 'fs-extra'; import { Response } from 'express'; import Path from 'path'; import _ from 'lodash'; import { corsHandler } from '../middleware/cors'; import { stringBody } from '../middleware/stringBody'; import { KapetaBodyRequest } from '../types'; import { HTMLPageEncoding, StormCodegenRequest, StormContextRequest, StormCreateBlockRequest, StormStream, } from './stream'; import { ConversationIdHeader, UIPagePrompt, UIPageEditRequest, BasePromptRequest, UIPageVoteRequest, UIPageGetVoteRequest, StormClient, } from './stormClient'; import { StormEvent, StormEventPage, StormEventPhaseType, UserJourneyScreen } from './events'; import { createPhaseEndEvent, createPhaseStartEvent, resolveOptions, StormDefinitions, StormEventParser, } from './event-parser'; import { StormCodegen } from './codegen'; import { assetManager } from '../assetManager'; import uuid from 'node-uuid'; import { getSystemBaseDir, getSystemBaseImplDir, SystemIdHeader, writeAssetToDisk, writePageToDisk, } from './page-utils'; import { UIServer } from './UIServer'; import { randomUUID } from 'crypto'; import { copyDirectory, createFuture, readFilesAndContent } from './utils'; import { getRemoteUrl } from '../utils/utils'; import stormService from '../stormService'; import { PageQueue } from './PageGenerator'; const UI_SERVERS: { [key: string]: UIServer } = {}; const router = Router(); router.use('/', corsHandler); router.use('/', stringBody); const samplesBaseDir = Path.join(__dirname, 'samples'); function convertPageEvent(screenData: StormEvent, innerConversationId: string, mainConversationId: string): StormEvent { if (screenData.type === 'PAGE') { const server: UIServer | undefined = UI_SERVERS[mainConversationId]; if (!server) { console.warn('No server found for conversation', mainConversationId); } screenData.payload.conversationId = innerConversationId; return { type: 'PAGE_URL', reason: screenData.reason, created: screenData.created, payload: { id: uuid.v4(), name: screenData.payload.name, title: screenData.payload.title, filename: screenData.payload.filename, description: screenData.payload.description, prompt: screenData.payload.prompt, path: screenData.payload.path, url: screenData.payload.content ? screenData.payload.path : '', method: screenData.payload.method, conversationId: innerConversationId, }, }; } return screenData; } router.post('/ui/serve/:systemId', async (req: KapetaBodyRequest, res: Response) => { const systemId = req.params.systemId as string | undefined; if (!systemId) { res.status(404).send({ error: 'Missing "systemId" in URL' }); return; } const svr = (UI_SERVERS[systemId] = UI_SERVERS[systemId] || new UIServer(systemId)); if (!svr.isRunning()) { await UI_SERVERS[systemId].start(); } res.status(200).send({ status: 'running', url: svr.getUrl(), resetUrl: svr.resolveUrlFromPath('/_reset') }); }); router.get('/ui/conversations', async (req: KapetaBodyRequest, res: Response) => { const local = await stormService.listLocalConversations(); const remote = await stormService.listRemoteConversations(); res.send({ local, remote, }); }); router.get('/ui/conversations/:systemId', async (req: KapetaBodyRequest, res: Response) => { const systemId = req.params.systemId as string; const eventsFile = getSystemBaseDir(systemId) + '/events.ndjson'; res.set('Content-Type', 'application/x-ndjson'); res.set('Access-Control-Expose-Headers', ConversationIdHeader); res.set(ConversationIdHeader, systemId); if (!FS.existsSync(eventsFile)) { res.status(404).send({ error: 'No events found' }); return; } res.send(FS.readFileSync(eventsFile)); }); router.put('/ui/conversations/:systemId', async (req: KapetaBodyRequest, res: Response) => { const systemId = req.params.systemId as string; const events: StormEvent[] = req.stringBody ? JSON.parse(req.stringBody) : []; await stormService.saveConversation(systemId, events); res.send({ ok: true }); }); router.delete('/ui/conversations/:systemId', async (req: KapetaBodyRequest, res: Response) => { const systemId = req.params.systemId as string; await stormService.deleteConversation(systemId); res.send({ ok: true }); }); router.post('/ui/conversations/:systemId/append', async (req: KapetaBodyRequest, res: Response) => { const systemId = req.params.systemId as string; const events: StormEvent[] = req.stringBody ? JSON.parse(req.stringBody) : []; await stormService.appendConversation(systemId, events); res.send({ ok: true }); }); router.post('/ui/create-system/:handle/:systemId', async (req: KapetaBodyRequest, res: Response) => { const systemId = req.params.systemId as string; const handle = req.params.handle as string; const srcDir = getSystemBaseDir(systemId); const destDir = getSystemBaseImplDir(systemId); res.set('Content-Type', 'application/x-ndjson'); res.set('Access-Control-Expose-Headers', ConversationIdHeader); res.set(ConversationIdHeader, systemId); sendEvent(res, createPhaseStartEvent(StormEventPhaseType.IMPLEMENT_APIS)); const pagesFromDisk = readFilesAndContent(srcDir); const client = new StormClient(handle, systemId); const pagesWithImplementation = await client.replaceMockWithAPICall({ pages: pagesFromDisk, systemId: systemId, }); await copyDirectory(srcDir, destDir, (fileName, content) => { // find the page from result1 and write the content to the file const page = pagesWithImplementation.find((p) => p.fileName === fileName); return page ? page.content : content; }); sendEvent(res, createPhaseEndEvent(StormEventPhaseType.IMPLEMENT_APIS)); sendEvent(res, createPhaseStartEvent(StormEventPhaseType.COMPOSE_SYSTEM_PROMPT)); // get the content of the pages const pageContents = pagesWithImplementation.map((page) => { return page.content; }); const prompt = await client.generatePrompt(pageContents); sendEvent(res, createPhaseEndEvent(StormEventPhaseType.COMPOSE_SYSTEM_PROMPT)); req.query.systemId = systemId; const promptRequest: BasePromptRequest = { prompt: prompt, skipImprovement: true, }; req.stringBody = JSON.stringify(promptRequest); await handleAll(req, res); }); router.post('/ui/create-system-simple/:handle/:systemId', async (req: KapetaBodyRequest, res: Response) => { const handle = req.params.handle as string; const systemId = req.params.systemId as string; const srcDir = getSystemBaseDir(systemId); //res.set('Content-Type', 'application/x-ndjson'); //res.set('Access-Control-Expose-Headers', ConversationIdHeader); //res.set(ConversationIdHeader, systemId); //sendEvent(res, createPhaseStartEvent(StormEventPhaseType.IMPLEMENT_APIS)); const client = new StormClient(handle, systemId); try { const pagesFromDisk = readFilesAndContent(srcDir); const pagesWithImplementation = await client.replaceMockWithAPICall({ pages: pagesFromDisk, systemId: systemId, }); //sendEvent(res, createPhaseEndEvent(StormEventPhaseType.IMPLEMENT_APIS)); //sendEvent(res, createPhaseStartEvent(StormEventPhaseType.COMPOSE_SYSTEM)); const allFiles = readFilesAndContent(srcDir, false).map((page) => { if (page.encoding == HTMLPageEncoding.TEXT) { const matchingFile = pagesWithImplementation.find( (pageWithImpl) => pageWithImpl.fileName === page.fileName ); if (matchingFile) { return matchingFile; } } return page; }); const systemUrl = await client.createSimpleBackend(handle, systemId, { pages: allFiles }); //sendEvent(res, {type: 'SYSTEM_READY', created: Math.floor(Date.now() / 1000), reason: 'System Ready', payload: { systemUrl: systemUrl }}); //sendEvent(res, createPhaseEndEvent(StormEventPhaseType.COMPOSE_SYSTEM)); //sendDone(res); res.json({ url: systemUrl }); } catch (err: any) { res.status(500).json({ error: err.message }); } finally { if (!res.closed) { res.end(); } } }); router.post('/ui/systems/:handle/:systemId/download', async (req: KapetaBodyRequest, res: Response) => { const systemId = req.params.systemId as string; const handle = req.params.handle as string; await stormService.installProjectById(handle, systemId); res.send({ ok: true }); }); router.post('/ui/systems/:handle/:systemId/upload', async (req: KapetaBodyRequest, res: Response) => { const systemId = req.params.systemId as string; const handle = req.params.handle as string; await stormService.uploadConversation(handle, systemId); res.send({ ok: true }); }); router.put('/ui/systems/:handle/:systemId/thumbnail', async (req: KapetaBodyRequest, res: Response) => { const systemId = req.params.systemId as string; await stormService.saveThumbnail(systemId, req.body as Buffer); res.send({ ok: true }); }); router.get('/ui/systems/:handle/:systemId/thumbnail.png', async (req: KapetaBodyRequest, res: Response) => { const systemId = req.params.systemId as string; const thumbnail = await stormService.getThumbnail(systemId); if (thumbnail) { res.set('Content-Type', 'image/png'); res.send(thumbnail); } else { res.status(404).send({ error: 'No thumbnail found' }); } }); router.delete('/ui/serve/:systemId', async (req: KapetaBodyRequest, res: Response) => { const systemId = req.params.systemId as string | undefined; if (!systemId) { res.status(404).send({ error: 'Missing "systemId" in URL' }); return; } const server = UI_SERVERS[systemId]; if (server) { server.close(); delete UI_SERVERS[systemId]; } res.status(200).json({ status: 'ok' }); }); /** * Edit a single page */ router.post('/:handle/ui/screen', async (req: KapetaBodyRequest, res: Response) => { try { const handle = req.params.handle as string; const conversationId = req.headers[ConversationIdHeader.toLowerCase()] as string | undefined; const systemId = req.headers[SystemIdHeader.toLowerCase()] as string | undefined; const aiRequest: UIPagePrompt = JSON.parse(req.stringBody ?? '{}'); aiRequest.storage_prefix = systemId ? systemId + '_' : 'mock_'; res.set('Content-Type', 'application/x-ndjson'); res.set('Access-Control-Expose-Headers', ConversationIdHeader); res.set(ConversationIdHeader, conversationId); const parentConversationId = systemId ?? ''; const queue = new PageQueue(handle, parentConversationId, '', 5); onRequestAborted(req, res, () => { queue.cancel(); }); const promises: Promise[] = []; queue.on('page', (data) => (systemId ? sendPageEvent(systemId, data, res) : undefined)); queue.on('error', (err) => { console.error('Failed to process page', err); sendError(err as any, res); }); queue.on('event', (event: StormEvent) => { if (event.type === 'FILE_CHUNK') { return; } sendEvent(res, event); }); await queue.addPrompt(aiRequest, conversationId, true); await queue.wait(); await Promise.allSettled(promises); sendDone(res); } catch (err: any) { sendError(err, res); } finally { if (!res.closed) { res.end(); } } }); router.post('/:handle/ui/iterative', async (req: KapetaBodyRequest, res: Response) => { const handle = req.params.handle as string; try { const conversationId = req.headers[ConversationIdHeader.toLowerCase()] as string | undefined; const aiRequest: BasePromptRequest = JSON.parse(req.stringBody ?? '{}'); const client = new StormClient(handle, conversationId); //todo is this correct we are using the landing page getConversationId down below as well const landingPagesStream = await client.createUILandingPages(aiRequest, conversationId); onRequestAborted(req, res, () => { landingPagesStream.abort(); }); res.set('Content-Type', 'application/x-ndjson'); res.set('Access-Control-Expose-Headers', ConversationIdHeader); res.set(ConversationIdHeader, landingPagesStream.getConversationId()); const promises: { [key: string]: Promise } = {}; const pageEventPromises: Promise[] = []; const systemId = landingPagesStream.getConversationId(); const systemPrompt = createFuture(); if (aiRequest.skipImprovement) { systemPrompt.resolve(aiRequest.prompt); } landingPagesStream.on('data', async (data: StormEvent) => { try { sendEvent(res, data); if (data.type === 'PROMPT_IMPROVE') { systemPrompt.resolve(data.payload.prompt); } if (data.type !== 'LANDING_PAGE') { return; } if (landingPagesStream.isAborted()) { return; } const landingPage = data.payload; if (landingPage.name in promises) { return; } // We add the landing pages to the prompt queue. // These will then be analysed - creating further pages as needed promises[landingPage.name] = pageQueue.addPrompt({ prompt: landingPage.create_prompt, method: 'GET', path: landingPage.path, description: landingPage.create_prompt, name: landingPage.name, title: landingPage.title, filename: landingPage.filename, storage_prefix: systemId + '_', // TODO: Add themes to this request type theme: '', }); } catch (e) { console.error('Failed to process event', e); } }); UI_SERVERS[systemId] = new UIServer(systemId); await UI_SERVERS[systemId].start(); waitForStormStream(landingPagesStream).then(() => { systemPrompt.resolve(aiRequest.prompt); }); const pageQueue = new PageQueue(handle, systemId, await systemPrompt.promise, 5); onRequestAborted(req, res, () => { pageQueue.cancel(); }); pageQueue.on('page', (screenData: StormEventPage) => sendPageEvent(landingPagesStream.getConversationId(), screenData, res) ); pageQueue.on('event', (event: StormEvent) => { if (event.type === 'FILE_CHUNK') { return; } sendEvent(res, event); }); pageQueue.on('error', (err) => { console.error('Failed to process page', err); sendError(err as any, res); }); await waitForStormStream(landingPagesStream); await pageQueue.wait(); await Promise.allSettled(pageEventPromises); if (landingPagesStream.isAborted()) { return; } sendDone(res); } catch (err) { sendError(err as Error, res); if (!res.closed) { res.end(); } } }); router.post('/:handle/ui', async (req: KapetaBodyRequest, res: Response) => { const handle = req.params.handle as string; try { const outerConversationId = (req.headers[ConversationIdHeader.toLowerCase()] as string | undefined) || randomUUID(); const aiRequest: BasePromptRequest = JSON.parse(req.stringBody ?? '{}'); const stormClient = new StormClient(handle, outerConversationId); // Get user journeys const userJourneysStream = await stormClient.createUIUserJourneys(aiRequest, outerConversationId); onRequestAborted(req, res, () => { userJourneysStream.abort(); }); res.set('Content-Type', 'application/x-ndjson'); res.set('Access-Control-Expose-Headers', ConversationIdHeader); res.set(ConversationIdHeader, outerConversationId); const uniqueUserJourneyScreens: Record = {}; let systemPrompt = aiRequest.prompt; userJourneysStream.on('data', (data: StormEvent) => { try { if (data.type === 'PROMPT_IMPROVE') { systemPrompt = data.payload.prompt; } if (data.type !== 'USER_JOURNEY') { sendEvent(res, data); return; } if (userJourneysStream.isAborted()) { return; } data.payload.screens.forEach((screen) => { if (!uniqueUserJourneyScreens[screen.name]) { screen.conversationId = randomUUID(); uniqueUserJourneyScreens[screen.name] = screen; } }); sendEvent(res, data); } catch (e) { console.error('Failed to process event', e); } }); userJourneysStream.on('error', (error) => { console.error('Error on userJourneysStream', error); userJourneysStream.abort(); sendError(error, res); }); let theme = ''; try { const themeStream = await stormClient.createTheme(aiRequest, outerConversationId); onRequestAborted(req, res, () => { themeStream.abort(); }); themeStream.on('data', (evt) => { if (evt.type === 'FILE_DONE') { theme = evt.payload.content; writeAssetToDisk(outerConversationId, evt).catch((err) => { sendEvent(res, { type: 'ERROR_INTERNAL', created: new Date().getTime(), payload: { error: err.message }, reason: 'Failed to save theme', }); }); } }); themeStream.on('error', (error) => { console.error(error); sendEvent(res, { type: 'ERROR_INTERNAL', created: new Date().getTime(), payload: { error: error.message }, reason: 'Failed to create theme', }); }); await waitForStormStream(themeStream); } catch (e: any) { console.error('Failed to generate theme', e); sendEvent(res, { type: 'ERROR_INTERNAL', created: new Date().getTime(), payload: { error: e.message }, reason: 'Failed to create theme', }); } await waitForStormStream(userJourneysStream); if (req.socket.closed) { return; } // Get the UI shells const shellsStream = await stormClient.createUIShells( { theme: theme || undefined, pages: Object.values(uniqueUserJourneyScreens).map((screen) => ({ name: screen.name, title: screen.title, filename: screen.filename, path: screen.path, method: screen.method, requirements: screen.requirements, })), }, outerConversationId ); onRequestAborted(req, res, () => { shellsStream.abort(); }); const queue = new PageQueue(handle, outerConversationId, systemPrompt, 5); queue.setUiTheme(theme); shellsStream.on('data', (data: StormEvent) => { //console.log('Processing shell event', data); sendEvent(res, data); if (data.type !== 'UI_SHELL') { return; } if (shellsStream.isAborted()) { return; } queue.addUiShell(data.payload); }); shellsStream.on('error', (error) => { console.error('Error on shellsStream', error); shellsStream.abort(); sendError(error, res); }); await waitForStormStream(shellsStream); if (req.socket.closed) { return; } UI_SERVERS[outerConversationId] = new UIServer(outerConversationId); await UI_SERVERS[outerConversationId].start(); sendEvent(res, { type: 'UI_SERVER_STARTED', reason: '', payload: { conversationId: outerConversationId, resetUrl: UI_SERVERS[outerConversationId].resolveUrlFromPath('/_reset'), }, created: Date.now(), }); onRequestAborted(req, res, () => { queue.cancel(); }); queue.on('page', (pageEvent: StormEventPage) => sendPageEvent(outerConversationId, pageEvent, res)); queue.on('event', (event: StormEvent) => { if (event.type === 'FILE_CHUNK') { return; } sendEvent(res, event); }); queue.on('error', (err) => { console.error('Failed to process page', err); sendError(err as any, res); }); for (const screen of Object.values(uniqueUserJourneyScreens)) { queue .addPrompt( { prompt: screen.requirements, method: screen.method, path: screen.path, description: screen.requirements, name: screen.name, title: screen.title, filename: screen.filename, storage_prefix: outerConversationId + '_', theme, }, screen.conversationId ) .catch((e) => { console.error('Failed to generate page for screen %s', screen.name, e); sendError(e as any, res); }); } if (userJourneysStream.isAborted()) { return; } await queue.wait(); sendDone(res); } catch (err) { sendError(err as Error, res); } finally { if (!res.closed) { res.end(); } } }); /** * Edit all pages */ router.post('/:handle/ui/edit', async (req: KapetaBodyRequest, res: Response) => { try { const handle = req.params.handle as string; const systemId = (req.headers[SystemIdHeader.toLowerCase()] || req.headers[ConversationIdHeader.toLowerCase()]) as string | undefined; const aiRequest: StormContextRequest = JSON.parse(req.stringBody ?? '{}'); const storagePrefix = systemId ? systemId + '_' : 'mock_'; const queue = new PageQueue(handle, systemId!, '', 5); onRequestAborted(req, res, () => { queue.cancel(); }); const promises: Promise[] = []; queue.on('page', (data) => { if (systemId) { const promise = sendPageEvent(systemId, data, res); promises.push(promise); return promise; } }); queue.on('event', (event) => { if (event.type === 'FILE_CHUNK') { return; } sendEvent(res, event); }); queue.on('error', (err) => { console.error('Failed to process page', err); sendError(err as any, res); }); const pages = aiRequest.prompt.pages.filter((page) => page.conversationId); if (pages.length === 0) { console.log('No pages to update', aiRequest.prompt.pages); sendDone(res); return; } await Promise.allSettled( pages.map((page) => { if (page.conversationId) { return queue.addPrompt( { title: page.title, name: page.name, method: page.method, path: page.path, description: page.description ?? '', filename: page.filename, prompt: aiRequest.prompt.prompt.prompt, storage_prefix: storagePrefix, }, page.conversationId, true, true // this is a global edit ); } }) ); await queue.wait(); await Promise.all(promises); sendDone(res); } catch (err: any) { sendError(err as Error, res); } finally { if (!res.closed) { res.end(); } } }); router.post('/ui/vote', async (req: KapetaBodyRequest, res: Response) => { const conversationId = (req.headers[ConversationIdHeader.toLowerCase()] as string | undefined) || ''; const aiRequest: UIPageVoteRequest = JSON.parse(req.stringBody ?? '{}'); const { topic, vote, mainConversationId } = aiRequest; try { const stormClient = new StormClient('', mainConversationId); await stormClient.voteUIPage(topic, conversationId, vote, mainConversationId); } catch (e: any) { res.status(500).send({ error: e.message }); } }); router.post('/ui/get-vote', async (req: KapetaBodyRequest, res: Response) => { const conversationId = (req.headers[ConversationIdHeader.toLowerCase()] as string | undefined) || ''; const aiRequest: UIPageGetVoteRequest = JSON.parse(req.stringBody ?? '{}'); const { topic, mainConversationId } = aiRequest; try { const stormClient = new StormClient('', mainConversationId); const vote = await stormClient.getVoteUIPage(topic, conversationId, mainConversationId); res.send({ vote }); } catch (e: any) { res.status(500).send({ error: e.message }); } }); router.post('/:handle/all', async (req: KapetaBodyRequest, res: Response) => { await handleAll(req, res); }); router.get('/conversations/:agentName/:systemId/redirect', async (req: KapetaBodyRequest, res: Response) => { const aiService = getRemoteUrl('ai-service', 'https://ai.kapeta.com'); const agentName = req.params.agentName as string; const systemId = req.params.systemId as string; res.redirect(aiService + `/v2/conversations/${encodeURIComponent(agentName)}/${encodeURIComponent(systemId)}`); }); async function handleAll(req: KapetaBodyRequest, res: Response) { const handle = req.params.handle as string; const systemId = (req.query.systemId as string) ?? undefined; try { const stormOptions = { ...(await resolveOptions()), systemId: systemId }; const eventParser = new StormEventParser(stormOptions); const conversationId = req.headers[ConversationIdHeader.toLowerCase()] as string | undefined; const aiRequest: BasePromptRequest = JSON.parse(req.stringBody ?? '{}'); const stormClient = new StormClient(handle, systemId); const metaStream = await stormClient.createMetadata(aiRequest, conversationId); onRequestAborted(req, res, () => { metaStream.abort(); }); // We check if the headers have been sent, because we might have already sent some data // before this function is called if (!res.headersSent) { res.set('Content-Type', 'application/x-ndjson'); res.set('Access-Control-Expose-Headers', ConversationIdHeader); res.set(ConversationIdHeader, metaStream.getConversationId()); } let currentPhase = StormEventPhaseType.META; // Helper to avoid sending the plan multiple times in a row const sendUpdatedPlan = _.debounce(sendDefinitions, 50, { maxWait: 200 }); metaStream.on('data', async (data: StormEvent) => { try { const result = await eventParser.processEvent(handle, data); switch (data.type) { case 'API_STREAM_START': case 'CREATE_API': case 'CREATE_MODEL': case 'CREATE_TYPE': if (currentPhase !== StormEventPhaseType.DEFINITIONS) { sendEvent(res, createPhaseEndEvent(StormEventPhaseType.META)); currentPhase = StormEventPhaseType.DEFINITIONS; sendEvent(res, createPhaseStartEvent(StormEventPhaseType.DEFINITIONS)); } break; } sendEvent(res, data); sendUpdatedPlan(res, result); } catch (e) { console.error('Failed to process event', e); } }); try { sendEvent(res, createPhaseStartEvent(StormEventPhaseType.META)); await waitForStormStream(metaStream); } finally { if (!metaStream.isAborted()) { sendEvent(res, createPhaseEndEvent(currentPhase)); } } if (metaStream.isAborted()) { return; } if (!eventParser.isValid()) { // We can't continue if the meta stream is invalid sendEvent(res, { type: 'ERROR_INTERNAL', payload: { error: eventParser.getError() }, reason: 'Failed to generate system', created: Date.now(), }); res.end(); return; } const result = await eventParser.toResult(handle, true); if (metaStream.isAborted()) { return; } // Cancel debounce, we don't need to send the plan again sendUpdatedPlan.cancel(); sendDefinitions(res, result); if (!req.query.skipCodegen) { try { sendEvent(res, createPhaseStartEvent(StormEventPhaseType.IMPLEMENTATION)); const stormCodegen = new StormCodegen( metaStream.getConversationId(), aiRequest.prompt, result.blocks, eventParser.getEvents(), systemId ); onRequestAborted(req, res, () => { stormCodegen.abort(); }); const codegenPromise = streamStormPartialResponse(stormCodegen.getStream(), res); await stormCodegen.process(); await codegenPromise; } finally { if (!metaStream.isAborted()) { sendEvent(res, createPhaseEndEvent(StormEventPhaseType.IMPLEMENTATION)); } } } sendDone(res); } catch (err: any) { sendError(err, res); if (!res.closed) { res.end(); } } } router.post('/block/create', async (req: KapetaBodyRequest, res: Response) => { const createRequest: StormCreateBlockRequest = JSON.parse(req.stringBody ?? '{}'); try { const ymlPath = Path.join(createRequest.newPath, 'kapeta.yml'); const isExisting = await FS.pathExists(createRequest.tmpPath); // The tmp folder might have a different name, so we need to update the asset after moving it over if (isExisting) { await FS.remove(`${createRequest.tmpPath}/kapeta.yml`); await FS.move(createRequest.tmpPath, createRequest.newPath, { overwrite: true, }); } // Create asset without running codegen if the asset already exists const shouldCodegen = !isExisting; const [asset] = await assetManager.createAsset(ymlPath, createRequest.definition, 'user', shouldCodegen); res.send(asset); } catch (err: any) { res.status(500).send({ error: err.message }); } }); router.post('/block/codegen', async (req: KapetaBodyRequest, res: Response) => { const body: StormCodegenRequest = JSON.parse(req.stringBody ?? '{}'); const conversationId = req.headers[ConversationIdHeader.toLowerCase()] as string | undefined; try { const stormCodegen = new StormCodegen(conversationId ?? '', body.prompt, [body.block], body.events || []); stormCodegen.setTmpDir(body.outDir); onRequestAborted(req, res, () => { stormCodegen.abort(); }); const codegenPromise = streamStormPartialResponse(stormCodegen.getStream(), res); await stormCodegen.process(); await codegenPromise; sendDone(res); } catch (err: any) { console.error('Failed to generate code', err); res.status(500).send({ error: err.message }); } }); function sendDefinitions(res: Response, result: StormDefinitions) { sendEvent(res, { type: 'DEFINITION_CHANGE', payload: result, reason: 'Updates to definition', created: Date.now(), }); } function sendDone(res: Response) { if (res.closed) { return; } sendEvent(res, { type: 'DONE', created: Date.now(), }); res.end(); } function sendError(err: Error, res: Response) { if (res.closed) { return; } const errorPayload = { error: err.message, stack: err.stack, }; console.error('Failed to send prompt', err); if (res.headersSent) { sendEvent(res, { type: 'ERROR_INTERNAL', created: Date.now(), payload: errorPayload, reason: 'Failed while sending prompt', }); } else { res.status(400).send(errorPayload); } } function waitForStormStream(result: StormStream) { return result.waitForDone(); } function streamStormPartialResponse(result: StormStream, res: Response) { return new Promise((resolve) => { result.on('data', (data) => { switch (data.type) { // todo: temporarily (for demo purposes) disable error messages when codegen fails case 'ERROR_INTERNAL': console.log('Error internal', data); return; } sendEvent(res, data); }); resolve(result.waitForDone()); }); } function sendEvent(res: Response, evt: StormEvent) { if (res.closed) { return; } res.write(JSON.stringify(evt) + '\n'); } function onRequestAborted(req: KapetaBodyRequest, res: Response, onAborted: () => void) { req.socket.on('close', () => { onAborted(); }); } async function sendPageEvent(mainConversationId: string, data: StormEventPage, res: Response) { if (data.payload.content) { try { await writePageToDisk(mainConversationId, data); } catch (err) { console.error('Failed to write page to disk', err); } } sendEvent(res, convertPageEvent(data, data.payload.conversationId, mainConversationId)); } export default router;