41 lines
1.8 KiB
JavaScript
41 lines
1.8 KiB
JavaScript
import net from 'node:net';
|
|
import fs from 'node:fs';
|
|
import path from 'node:path';
|
|
import { mkdir, rm } from 'node:fs/promises';
|
|
|
|
const MAX_LINE = 64 * 1024;
|
|
|
|
export async function createEventSocket({ socketPath, onMessage }) {
|
|
await mkdir(path.dirname(socketPath), { recursive: true });
|
|
try { await rm(socketPath, { force: true }); } catch {}
|
|
const clients = new Set();
|
|
const server = net.createServer((socket) => {
|
|
clients.add(socket);
|
|
socket.on('error', () => {});
|
|
socket.once('close', () => clients.delete(socket));
|
|
const invalid = () => { if (!socket.destroyed) socket.write(JSON.stringify({ error: 'invalid IPC message' }) + '\n'); };
|
|
let buffer = '';
|
|
socket.setEncoding('utf8');
|
|
socket.on('data', (chunk) => {
|
|
buffer += chunk;
|
|
let index;
|
|
while ((index = buffer.indexOf('\n')) >= 0) {
|
|
const line = buffer.slice(0, index); buffer = buffer.slice(index + 1);
|
|
if (Buffer.byteLength(line) > MAX_LINE) { socket.destroy(new Error('IPC message too large')); return; }
|
|
try { Promise.resolve(onMessage?.(JSON.parse(line), socket)).catch(invalid); } catch { invalid(); }
|
|
}
|
|
if (Buffer.byteLength(buffer) > MAX_LINE) socket.destroy(new Error('IPC message too large'));
|
|
});
|
|
});
|
|
await new Promise((resolve, reject) => { server.once('error', reject); server.listen(socketPath, resolve); });
|
|
return { server, socketPath, close: () => new Promise((resolve) => { for (const client of clients) client.destroy(); server.close(() => fs.unlink(socketPath, () => resolve())); }) };
|
|
}
|
|
|
|
export function sendEvent(socket, event, payload = {}) {
|
|
const message = JSON.stringify({ event, ...payload });
|
|
if (Buffer.byteLength(message) > MAX_LINE) throw new Error('IPC event too large; use a file path for bulk data');
|
|
socket.write(`${message}\n`);
|
|
}
|
|
|
|
export { MAX_LINE };
|