Files
stack/packages/bus/src/server.mjs
T
jason.woltjeandClaude Opus 5.5 7e73c2cd13 feat(tasks): the Vikunja v2 adapter, broker task verbs and sync (row 38, S3, darkwing)
packages/tasks adds the Vikunja v2 client, the eight task verbs, the
board-plus-cursor poll with its 60 s window and the digest. The broker
gains the task verbs and boots trackers from the boot config (lead
decisions 66 to 68). Due dates are truncated to the second and recorded
as truncated (B1). A write that lands but whose final read fails counts
as landed, in update and in create (B2).

Candidate agents/darkwing/work/slice1-s3, build-r2.patch 71ce87e6,
manifest e10e30e3 (28 files). Filbert approved round 2 on #1520
(comment 26853). Darkwing's post-reset rerun: test-release 14/14,
test-task 98/98 (comment 26857).

Integration gate in a worktree on c4baf779 with the patch applied:
bus 67, business 60, control-board 124, discord 173, ledger 78,
mosaic 69, queue 148, seat 19, tasks 51 and webui 14, all with no
failures. Conversation is 149/3. The three cohort kill cases (K1, K3,
K10) fail the same on the unpatched base, and the patch doesn't touch
the package. Every scripts/test-*.sh is green. test-release 14/14 and
test-task 98/98 ran on the existing gate2 compose network, because the
host's Docker address pools are exhausted. No network was created or
pruned.

Co-Authored-By: Claude Opus 5.5 <[email protected]>
2026-10-09 07:40:48 -05:00

108 lines
3.8 KiB
JavaScript

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;
}
},
};
}