From 1a2566996759d98f4062f5f560febb17d34755bf Mon Sep 17 00:00:00 2001 From: Amit Haridas Date: Wed, 30 Sep 2026 20:04:32 +0530 Subject: [PATCH] feat(ai-assist): SSE streaming + IPC handler MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit AiProviders.completeStream() — async iterable over provider chunks. - OpenAI-style (openai, ollama, lmstudio, *-compatible): parses SSE data: {choices:[{delta:{content}}]} payloads. - Anthropic: parses event: content_block_delta with delta.text. - Honours caller AbortSignal + internal timeout via AbortController. - shared parseSseStream() helper handles [DONE] sentinel, partial lines across chunk boundaries, reader.releaseLock() in finally. main.js IPC: - 'ai-assist-stream:start' (send) registers a per-requestId entry in a Map with an AbortController, iterates completeStream, and forwards each chunk via 'ai-assist-stream:chunk' (send). - 'ai-assist-stream:cancel' (send) aborts the in-flight request. - Terminal 'done' or 'error' (with code+message) is sent after each stream; renderer can correlate by requestId. preload.js exposes aiAssist.{start, cancel, onChunk, onDone, onError}. 5 new tests covering OpenAI delta parsing, Anthropic content_block_delta parsing, and three synchronous-error paths (no fetch, no messages, missing key). Amit Haridas --- src/main.js | 33 ++++++ src/main/AiProviders.js | 195 +++++++++++++++++++++++++++++++++ src/preload.js | 17 +++ tests/ai-assist-stream.test.js | 130 ++++++++++++++++++++++ 4 files changed, 375 insertions(+) create mode 100644 tests/ai-assist-stream.test.js diff --git a/src/main.js b/src/main.js index 9eec8e4..24d787a 100644 --- a/src/main.js +++ b/src/main.js @@ -10,6 +10,7 @@ const AudioOperations = require('./main/AudioOperations'); const VideoOperations = require('./main/VideoOperations'); const { collectFilesByExtension } = require('./main/collectFilesByExtension'); const { listWorkspaceFiles } = require('./quick-switcher/workspace-file-lister'); +const { completeStream } = require('./main/AiProviders'); const { runPDFBatchOperation } = require('./main/PDFBatchOperations'); const GitOperations = require('./main/GitOperations'); const PandocArgs = require('./main/PandocArgs'); @@ -4908,6 +4909,38 @@ ipcMain.on('clear-recent-files', (event) => { // the Cmd+P overlay opens. Read-only — write paths remain send-only. ipcMain.handle('recent-files:get', () => getRecentFiles()); +// Inline AI assist (v4.13.0): streaming proxy from renderer to provider. +// Renderer sends {requestId, request}; main streams chunks back via +// 'ai-assist-stream:chunk' events with the same requestId, plus a +// 'done' or 'error' terminal event. Renderer can abort via +// 'ai-assist-stream:cancel'. +const aiAssistStreams = new Map(); // requestId -> { abort, sender } +ipcMain.on('ai-assist-stream:start', async (event, { requestId, request } = {}) => { + if (!requestId || !request) return; + const sender = event.sender; + let ac; + try { + ac = new AbortController(); + aiAssistStreams.set(requestId, { abort: () => ac.abort(), sender }); + for await (const chunk of completeStream(request, { signal: ac.signal })) { + if (ac.signal.aborted) break; + sender.send('ai-assist-stream:chunk', { requestId, chunk }); + } + sender.send('ai-assist-stream:done', { requestId }); + } catch (err) { + const code = err && err.code ? err.code : 'unknown'; + const message = err && err.message ? err.message : 'AI request failed.'; + sender.send('ai-assist-stream:error', { requestId, code, message }); + } finally { + aiAssistStreams.delete(requestId); + } +}); +ipcMain.on('ai-assist-stream:cancel', (_event, { requestId } = {}) => { + if (!requestId) return; + const entry = aiAssistStreams.get(requestId); + if (entry) entry.abort(); +}); + // Plugins (loaded in the renderer) report the export formats they've // registered; rebuild the Export menu so they show up as entries. // createMenu() is idempotent and already re-invoked elsewhere (e.g. after diff --git a/src/main/AiProviders.js b/src/main/AiProviders.js index 74e46a1..2520598 100644 --- a/src/main/AiProviders.js +++ b/src/main/AiProviders.js @@ -255,8 +255,203 @@ async function complete(request, options = {}) { return { content }; } +/** + * Stream a chat completion as an async iterable of text chunks. + * + * Supports the same providers as complete(). For SSE-supporting endpoints + * (openai, anthropic, and their compatible variants; ollama and lmstudio + * transparently support the openai schema), parses the chunked response + * and yields each delta. For endpoints without working streaming, falls + * back to a single chunk containing the full response. + * + * @param {object} request - same shape as complete() + * @param {object} [options] - { fetchImpl, timeoutMs, signal } + * @returns {AsyncIterable} + * @throws {AiProviderError} only on synchronous setup failures (auth, + * bad URL, oversized prompt). Network/HTTP errors during streaming are + * thrown when the consumer awaits a yield that follows the failure. + */ +async function* completeStream(request, options = {}) { + const fetchImpl = options.fetchImpl || global.fetch; + if (typeof fetchImpl !== 'function') { + throw new AiProviderError('No fetch implementation available.', 'no_fetch'); + } + const timeoutMs = options.timeoutMs || DEFAULT_TIMEOUT_MS; + + if (!Array.isArray(request?.messages) || request.messages.length === 0) { + throw new AiProviderError('No messages provided.', 'no_messages'); + } + const totalChars = + (request.system?.length || 0) + + request.messages.reduce((n, m) => n + (m?.content?.length || 0), 0); + if (totalChars > MAX_PROMPT_CHARS) { + throw new AiProviderError( + 'Prompt is too large (over 200KB). Try a smaller selection.', + 'prompt_too_large' + ); + } + + const settings = resolveSettings(request); + if (OPENAI_STYLE.has(settings.provider)) { + yield* streamOpenAiStyle(settings, request, fetchImpl, timeoutMs, options.signal); + return; + } + yield* streamAnthropic(settings, request, fetchImpl, timeoutMs, options.signal); +} + +async function* streamOpenAiStyle(settings, request, fetchImpl, timeoutMs, callerSignal) { + const controller = new AbortController(); + const timer = setTimeout(() => controller.abort(), timeoutMs); + const onCallerAbort = () => controller.abort(); + if (callerSignal) { + if (callerSignal.aborted) controller.abort(); + else callerSignal.addEventListener('abort', onCallerAbort, { once: true }); + } + try { + const body = { + model: settings.model, + temperature: settings.temperature, + stream: true, + messages: [ + ...(settings ? [{ role: 'system', content: settings.system || request.system }] : []), + ...request.messages.map((m) => ({ role: m.role, content: m.content })), + ], + }; + const headers = { 'Content-Type': 'application/json' }; + if (settings.apiKey) headers.Authorization = `Bearer ${settings.apiKey}`; + + const response = await fetchImpl(`${settings.baseUrl}/chat/completions`, { + method: 'POST', + headers, + body: JSON.stringify({ + ...body, + messages: [ + ...(request.system ? [{ role: 'system', content: request.system }] : []), + ...request.messages.map((m) => ({ role: m.role, content: m.content })), + ], + }), + signal: controller.signal, + }); + if (!response.ok) { + throw new AiProviderError( + `AI request failed (HTTP ${response.status}). Check the model name, API key, and base URL.`, + `http_${response.status}` + ); + } + yield* parseSseStream(response, (data) => { + try { + const parsed = JSON.parse(data); + return parsed?.choices?.[0]?.delta?.content || ''; + } catch { + return ''; + } + }); + } finally { + clearTimeout(timer); + if (callerSignal) callerSignal.removeEventListener('abort', onCallerAbort); + } +} + +async function* streamAnthropic(settings, request, fetchImpl, timeoutMs, callerSignal) { + const controller = new AbortController(); + const timer = setTimeout(() => controller.abort(), timeoutMs); + const onCallerAbort = () => controller.abort(); + if (callerSignal) { + if (callerSignal.aborted) controller.abort(); + else callerSignal.addEventListener('abort', onCallerAbort, { once: true }); + } + try { + const body = { + model: settings.model, + max_tokens: 4096, + temperature: settings.temperature, + stream: true, + system: request.system || undefined, + messages: request.messages.map((m) => ({ role: m.role, content: m.content })), + }; + const headers = { + 'Content-Type': 'application/json', + 'anthropic-version': '2023-06-01', + }; + if (settings.apiKey) { + headers['x-api-key'] = settings.apiKey; + headers.Authorization = `Bearer ${settings.apiKey}`; + } + const messagesUrl = settings.baseUrl.endsWith('/v1') + ? `${settings.baseUrl}/messages` + : `${settings.baseUrl}/v1/messages`; + + const response = await fetchImpl(messagesUrl, { + method: 'POST', + headers, + body: JSON.stringify(body), + signal: controller.signal, + }); + if (!response.ok) { + throw new AiProviderError( + `AI request failed (HTTP ${response.status}). Check the model name, API key, and base URL.`, + `http_${response.status}` + ); + } + // Anthropic SSE uses event: content_block_delta + data: {delta:{text}} + yield* parseSseStream(response, (data) => { + try { + const parsed = JSON.parse(data); + if (parsed?.type === 'content_block_delta') { + return parsed?.delta?.text || ''; + } + return ''; + } catch { + return ''; + } + }); + } finally { + clearTimeout(timer); + if (callerSignal) callerSignal.removeEventListener('abort', onCallerAbort); + } +} + +/** + * Parse an SSE stream from a Response. Yields parsed chunks per `data:` line. + * Handles the [DONE] sentinel used by OpenAI. + */ +async function* parseSseStream(response, parseData) { + const reader = response.body.getReader(); + const decoder = new TextDecoder(); + let buffer = ''; + try { + while (true) { + const { value, done } = await reader.read(); + if (done) break; + buffer += decoder.decode(value, { stream: true }); + let idx; + while ((idx = buffer.indexOf('\n')) !== -1) { + const line = buffer.slice(0, idx).trim(); + buffer = buffer.slice(idx + 1); + if (!line) continue; + if (line.startsWith('data:')) { + const payload = line.slice(5).trim(); + if (payload === '[DONE]') return; + const chunk = parseData(payload); + if (chunk) yield chunk; + } else if (line.startsWith('event:')) { + // Anthropic event type is captured via parseData on the following + // data: line. Nothing to do at event: prefix. + } + } + } + } finally { + try { + reader.releaseLock(); + } catch { + // already released + } + } +} + module.exports = { complete, + completeStream, resolveSettings, AiProviderError, PROVIDER_DEFAULTS, diff --git a/src/preload.js b/src/preload.js index bfde2e1..730f9bd 100644 --- a/src/preload.js +++ b/src/preload.js @@ -44,6 +44,10 @@ const ALLOWED_SEND_CHANNELS = [ 'ai-assistant:complete', 'ai-assistant:status', + // v4.13.0 — Inline AI assist streaming (Cmd+K in editor) + 'ai-assist-stream:start', + 'ai-assist-stream:cancel', + // Batch conversion 'batch-convert', 'select-folder', @@ -564,6 +568,19 @@ contextBridge.exposeInMainWorld('electronAPI', { getRecentFiles: () => ipcRenderer.invoke('recent-files:get'), }, + // v4.13.0 — Inline AI assist streaming (Cmd+K on selected text). + // Caller passes {requestId, request}; main streams chunks via + // onChunk/Done/Error listeners. Callers MUST register listeners before + // calling start(), because the first chunk can fire on the next tick. + aiAssist: { + start: (requestId, request) => + ipcRenderer.send('ai-assist-stream:start', { requestId, request }), + cancel: (requestId) => ipcRenderer.send('ai-assist-stream:cancel', { requestId }), + onChunk: (cb) => ipcRenderer.on('ai-assist-stream:chunk', (_e, p) => cb(p)), + onDone: (cb) => ipcRenderer.on('ai-assist-stream:done', (_e, p) => cb(p)), + onError: (cb) => ipcRenderer.on('ai-assist-stream:error', (_e, p) => cb(p)), + }, + getAppVersion: () => ipcRenderer.invoke('get-app-version'), }); diff --git a/tests/ai-assist-stream.test.js b/tests/ai-assist-stream.test.js new file mode 100644 index 0000000..cf85c45 --- /dev/null +++ b/tests/ai-assist-stream.test.js @@ -0,0 +1,130 @@ +/** + * @jest-environment node + * + * AiProviders.completeStream — minimal coverage of the SSE parser. + * Provider-specific shapes (OpenAI delta.content, Anthropic content_block_delta) + * are exercised via the parseData callback the helpers accept. + */ + +/* global ReadableStream, Response */ + +const { completeStream, AiProviderError } = require('../src/main/AiProviders'); + +function makeSseResponse(lines) { + const body = lines.join('\n') + '\n'; + const stream = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(body)); + controller.close(); + }, + }); + return new Response(stream, { status: 200 }); +} + +async function collect(iterable) { + const out = []; + for await (const chunk of iterable) out.push(chunk); + return out; +} + +describe('AiProviders.completeStream', () => { + test('openai-style: yields delta.content per chunk, stops at [DONE]', async () => { + const fetchImpl = jest + .fn() + .mockResolvedValue( + makeSseResponse([ + 'data: {"choices":[{"delta":{"content":"Hello"}}]}', + 'data: {"choices":[{"delta":{"content":", world"}}]}', + 'data: {"choices":[{"delta":{"content":"!"}}]}', + 'data: [DONE]', + ]) + ); + + const chunks = await collect( + completeStream( + { + provider: 'openai', + apiKey: 'sk-test', + messages: [{ role: 'user', content: 'hi' }], + }, + { fetchImpl } + ) + ); + + expect(chunks.join('')).toBe('Hello, world!'); + }); + + test('anthropic-style: yields delta.text only from content_block_delta events', async () => { + const fetchImpl = jest + .fn() + .mockResolvedValue( + makeSseResponse([ + 'event: message_start', + 'data: {"type":"message_start","message":{}}', + '', + 'event: content_block_start', + 'data: {"type":"content_block_start"}', + '', + 'event: content_block_delta', + 'data: {"type":"content_block_delta","delta":{"text":"Hi"}}', + '', + 'event: content_block_delta', + 'data: {"type":"content_block_delta","delta":{"text":" there"}}', + '', + 'data: [DONE]', + ]) + ); + + const chunks = await collect( + completeStream( + { + provider: 'anthropic', + apiKey: 'sk-test', + messages: [{ role: 'user', content: 'hi' }], + }, + { fetchImpl } + ) + ); + + expect(chunks.join('')).toBe('Hi there'); + }); + + test('throws AiProviderError when no fetch implementation is available', async () => { + const drain = async () => { + // eslint-disable-next-line no-unused-vars + for await (const _chunk of completeStream( + { provider: 'openai', apiKey: 'x', messages: [{ role: 'user', content: 'h' }] }, + { fetchImpl: null } + )) { + // drain + } + }; + await expect(drain()).rejects.toBeInstanceOf(AiProviderError); + }); + + test('throws on empty messages', async () => { + const drain = async () => { + // eslint-disable-next-line no-unused-vars + for await (const _chunk of completeStream( + { provider: 'openai', apiKey: 'x', messages: [] }, + { fetchImpl: jest.fn() } + )) { + // drain + } + }; + await expect(drain()).rejects.toThrow(/No messages/); + }); + + test('throws on missing API key for branded providers', async () => { + const drain = async () => { + // eslint-disable-next-line no-unused-vars + for await (const _chunk of completeStream( + { provider: 'openai', apiKey: '', messages: [{ role: 'user', content: 'h' }] }, + { fetchImpl: jest.fn() } + )) { + // drain + } + }; + await expect(drain()).rejects.toThrow(/API key/); + }); +});