diff --git a/frontend/app/api/launch_kit/route.js b/frontend/app/api/launch_kit/route.js index a6316604..f0748432 100644 --- a/frontend/app/api/launch_kit/route.js +++ b/frontend/app/api/launch_kit/route.js @@ -18,6 +18,72 @@ import { readGenerationRequestBody } from "../../../lib/server/generationRequest const OWNER_ONLY_ENDPOINT_PROVIDERS = new Set(["custom", "ollama", "lmstudio"]); +function safeGenerationFailure(error) { + if (error instanceof ProviderError) { + const providerError = providerErrorPayload(error); + return { + ok: false, + error: providerError.message, + providerError, + warnings: [providerError.message], + }; + } + return { + ok: false, + error: "SignalFlow could not complete campaign generation.", + warnings: ["Campaign generation failed unexpectedly. Retry deliberately or inspect server diagnostics with the correlation context."], + }; +} + +function streamGeneration({ generationInput, generationConfig, warnings }) { + const encoder = new TextEncoder(); + const stream = new ReadableStream({ + start(controller) { + let closed = false; + const write = (value) => { + if (closed) return; + try { + controller.enqueue(encoder.encode(`${JSON.stringify(value)}\n`)); + } catch { + closed = true; + } + }; + Promise.resolve().then(async () => { + try { + const result = await generateStudioPackage({ + ...generationInput, + config: { + ...generationConfig, + onProgress: (progress) => write({ type: "progress", progress }), + }, + }); + const allWarnings = Array.from(new Set([...warnings, ...(result.warnings || [])])); + write({ type: "result", data: { ...result, warnings: allWarnings } }); + } catch (error) { + write({ type: "error", data: safeGenerationFailure(error) }); + } finally { + if (!closed) { + closed = true; + try { + controller.close(); + } catch { + // The browser may have cancelled the stream already. + } + } + } + }); + }, + }); + + return new Response(stream, { + status: 200, + headers: { + "Content-Type": "application/x-ndjson; charset=utf-8", + "Cache-Control": "no-store", + }, + }); +} + export const maxDuration = 60; export async function POST(request) { @@ -214,7 +280,7 @@ export async function POST(request) { } void enableAutoCapture; - const result = await generateStudioPackage({ + const generationInput = { projectName, notes, audience, @@ -228,13 +294,22 @@ export async function POST(request) { model_name: providerModelName || modelName, model_endpoint: providerBaseUrl || modelEndpoint, appUrl, - config: { - apiKey: providerApiKey, - baseUrl: providerBaseUrl, - modelName: providerModelName, - allowServerKey: isOwner, - signal: request.signal, - }, + }; + const generationConfig = { + apiKey: providerApiKey, + baseUrl: providerBaseUrl, + modelName: providerModelName, + allowServerKey: isOwner, + signal: request.signal, + }; + + if (request.headers.get("accept")?.includes("application/x-ndjson")) { + return streamGeneration({ generationInput, generationConfig, warnings }); + } + + const result = await generateStudioPackage({ + ...generationInput, + config: generationConfig, }); const allWarnings = Array.from(new Set([...warnings, ...(result.warnings || [])])); @@ -243,25 +318,12 @@ export async function POST(request) { headers: { "Content-Type": "application/json" }, }); } catch (error) { - if (error instanceof ProviderError) { - const providerError = providerErrorPayload(error); - return new Response(JSON.stringify({ - ok: false, - error: providerError.message, - providerError, - warnings: [providerError.message], - }), { - status: providerError.httpStatus && providerError.httpStatus >= 400 ? providerError.httpStatus : 502, - headers: { "Content-Type": "application/json" }, - }); - } - - return new Response(JSON.stringify({ - ok: false, - error: "SignalFlow could not complete campaign generation.", - warnings: ["Campaign generation failed unexpectedly. Retry deliberately or inspect server diagnostics with the correlation context."], - }), { - status: 500, + const failure = safeGenerationFailure(error); + const status = error instanceof ProviderError + ? (error.httpStatus && error.httpStatus >= 400 ? error.httpStatus : 502) + : 500; + return new Response(JSON.stringify(failure), { + status, headers: { "Content-Type": "application/json" }, }); } diff --git a/frontend/app/page.js b/frontend/app/page.js index c62b8946..8eaf007d 100644 --- a/frontend/app/page.js +++ b/frontend/app/page.js @@ -207,6 +207,58 @@ async function readJsonResponse(response, fallbackMessage) { throw new Error(response.ok ? fallbackMessage : `${fallbackMessage} (HTTP ${response.status})`); } +async function readGenerationResponse(response, onProgress, fallbackMessage) { + const contentType = String(response.headers.get("content-type") || "").toLowerCase(); + if (!contentType.includes("application/x-ndjson") || !response.body) { + return readJsonResponse(response, fallbackMessage); + } + + const reader = response.body.getReader(); + const decoder = new TextDecoder(); + let buffer = ""; + let finalData = null; + + const processLine = (line) => { + if (!line.trim()) return; + const event = safeJsonParse(line, null); + if (!event || typeof event !== "object") throw new Error(fallbackMessage); + if (event.type === "progress" && event.progress) { + onProgress?.(event.progress); + return; + } + if ((event.type === "result" || event.type === "error") && event.data) { + finalData = event.data; + } + }; + + while (true) { + const { done, value } = await reader.read(); + if (done) break; + buffer += decoder.decode(value, { stream: true }); + const lines = buffer.split("\n"); + buffer = lines.pop() || ""; + for (const line of lines) processLine(line); + } + buffer += decoder.decode(); + if (buffer.trim()) processLine(buffer); + + if (finalData && typeof finalData === "object") return finalData; + throw new Error(fallbackMessage); +} + +function generationProgressLabel(status) { + const labels = { + queued: "Queued", + generating: "Generating", + revising: "Revising", + complete: "Complete", + needs_review: "Needs review", + failed: "Failed", + cancelled: "Cancelled", + }; + return labels[String(status || "")] || "Preparing"; +} + function downloadText(filename, value, type = "text/plain") { const blob = new Blob([value], { type }); const url = URL.createObjectURL(blob); @@ -394,6 +446,7 @@ export default function Home() { const [files, setFiles] = useState([]); const [documentText, setDocumentText] = useState([]); const [busy, setBusy] = useState(false); + const [generationProgress, setGenerationProgress] = useState(null); const [message, setMessage] = useState(null); const [strategyReview, setStrategyReview] = useState(null); const [library, setLibrary] = useState([]); @@ -729,6 +782,7 @@ const sourceAndChannelsReady = sourceSignals > 0 && channels.length > 0; setFiles([]); setDocumentText([]); setPublishOptions({ reddit: { subreddit: "", title: "" } }); + setGenerationProgress(null); setMessage(null); navigateSection("studio"); } @@ -899,7 +953,10 @@ ${extractedText}`); }); const response = await fetch("/api/launch_kit", { method: "POST", - headers: authHeaders({ "Content-Type": "application/json" }), + headers: authHeaders({ + "Content-Type": "application/json", + Accept: "application/x-ndjson", + }), signal, body: JSON.stringify({ project_name: form.projectName.trim() || "Untitled campaign", @@ -930,7 +987,11 @@ ${extractedText}`); }), }); - const data = await readJsonResponse(response, "SignalFlow returned an unreadable generation response."); + const data = await readGenerationResponse( + response, + setGenerationProgress, + "SignalFlow returned an unreadable generation response.", + ); if (data.code === "strategy_quality_blocked" && data.strategy_review) { return { strategyBlocked: true, data }; } @@ -950,6 +1011,7 @@ ${extractedText}`); } function beginGenerationRequest() { + setGenerationProgress(null); const controller = new AbortController(); generationAbortRef.current = controller; return controller; @@ -963,6 +1025,9 @@ ${extractedText}`); const controller = generationAbortRef.current; if (!controller || controller.signal.aborted) return; controller.abort(); + setGenerationProgress((previous) => previous + ? { ...previous, phase: "cancelled", status: "cancelled" } + : previous); setMessage({ type: "warning", text: "Cancelling generation. Existing drafts will remain unchanged." }); } @@ -2301,12 +2366,41 @@ async function exportZip() {