import { createServer } from 'node:net'; import { lstatSync, chmodSync, unlinkSync } from 'node:fs'; import { BusError, TASK_VERBS } from './broker.mjs'; const LIMIT = 65536; // One request per connection. There is deliberately no reconnect/retry or SQL verb. // `tasks` ({handle, timeout}) is the trusted S3 adapter; only the eight task verbs reach it. export async function serve({ broker, path, authenticateHuman = null, timeout = 5000, tasks = null }) { try { lstatSync(path); throw new BusError('socket-exists'); } catch (e) { if (e.code !== 'ENOENT') throw e; } const sockets = new Set(); const server = createServer({ allowHalfOpen: true }, (socket) => { sockets.add(socket); socket.on('close', () => sockets.delete(socket)); socket.on('error', () => {}); let input = Buffer.alloc(0), done = false; const finish = (result) => { if (done) return; done = true; socket.end(JSON.stringify(result) + '\n'); }; socket.setTimeout(timeout, () => { finish({ ok: false, error: 'request-timeout' }); socket.destroySoon(); }); socket.on('data', (chunk) => { if (done) return; if (input.length + chunk.length > LIMIT) { finish({ ok: false, error: 'request-too-large' }); return; } input = Buffer.concat([input, chunk]); const end = input.indexOf(10); if (end < 0) return; let transient; try { if (end !== input.length - 1) throw new BusError('invalid-envelope'); const r = JSON.parse(input.subarray(0, end).toString('utf8')); if ( !r || Array.isArray(r) || typeof r !== 'object' || Object.keys(r).some((k) => !['cap', 'human', 'verb', 'args'].includes(k)) || typeof r.verb !== 'string' || Boolean(r.cap) === Boolean(r.human) ) throw new BusError('invalid-envelope'); let cap = r.cap; if (r.human) { if (!authenticateHuman) throw new BusError('human-required'); transient = cap = authenticateHuman(r.human); } if (typeof cap !== 'string') throw new BusError('unauthenticated'); if (tasks && TASK_VERBS.includes(r.verb)) { const own = transient; transient = null; // A Vikunja write can outlast the idle timeout; a late reply is still outcome-unknown to the client. socket.setTimeout(tasks.timeout); broker .requestTask(cap, { verb: r.verb, args: r.args ?? {} }, tasks.handle) .then( (result) => finish({ ok: true, result }), (e) => finish({ ok: false, error: e instanceof BusError ? e.code : 'adapter-failed' }), ) .finally(() => { if (own) broker.disconnect(own); }); return; } const result = broker.request(cap, { verb: r.verb, args: r.args ?? {} }); finish({ ok: true, result }); } catch (e) { finish({ ok: false, error: e instanceof BusError ? e.code : 'invalid-envelope' }); } finally { if (transient) broker.disconnect(transient); } }); socket.on('end', () => { if (!done) finish({ ok: false, error: 'incomplete-request' }); }); }); await new Promise((resolve, reject) => { server.once('error', reject); server.listen(path, () => { server.removeListener('error', reject); resolve(); }); }); chmodSync(path, 0o600); const owned = lstatSync(path); return { async close() { for (const s of sockets) s.destroy(); await new Promise((resolve) => server.close(resolve)); try { const cur = lstatSync(path); if (cur.ino === owned.ino && cur.dev === owned.dev) unlinkSync(path); } catch (e) { if (e.code !== 'ENOENT') throw e; } }, }; }