mirror of
https://github.com/amitwh/markdown-converter.git
synced 2026-10-01 09:19:34 +05:30
feat(ai-assist): SSE streaming + IPC handler
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
This commit is contained in:
+33
@@ -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
|
||||
|
||||
@@ -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<string>}
|
||||
* @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,
|
||||
|
||||
@@ -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'),
|
||||
});
|
||||
|
||||
|
||||
@@ -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/);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user