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]>
205 lines
8.6 KiB
JavaScript
205 lines
8.6 KiB
JavaScript
import { lstatSync } from 'node:fs';
|
|
import { BusError } from '../../bus/src/broker.mjs';
|
|
import { client, apiRoot } from './vikunja.mjs';
|
|
import { startup } from './startup.mjs';
|
|
import { tick, reconcile } from './sync.mjs';
|
|
import { verbs } from './verbs.mjs';
|
|
// S3's adapter for startBroker({tasks}). `trackers` is the host's plain-data boot config:
|
|
// {<business>: {baseUrl, project, pollSeconds, reconcileMinutes}}. A business with no entry has
|
|
// no tracker, and its task verbs refuse. Startup runs in the background so the broker boots
|
|
// without waiting on Vikunja; verbs refuse with tracker-starting until it passes.
|
|
export const TIMEOUT = 60000;
|
|
// Refusals that a later startup attempt can clear. A scope or install refusal waits for a restart.
|
|
const TRANSIENT = new Set(['tracker-unavailable', 'tracker-rate-limited', 'service-failed']);
|
|
const code = (e) => (e instanceof BusError ? e.code : 'adapter-failed');
|
|
const positive = (x) => Number.isSafeInteger(x) && x > 0;
|
|
function settings(name, t, b) {
|
|
if (!b || !t || typeof t !== 'object') throw new BusError('tracker-config');
|
|
apiRoot(t.baseUrl);
|
|
const pollSeconds = t.pollSeconds ?? 30,
|
|
reconcileMinutes = t.reconcileMinutes ?? 60;
|
|
if (!positive(t.project) || !Number.isSafeInteger(pollSeconds) || pollSeconds < 10 || !positive(reconcileMinutes))
|
|
throw new BusError('tracker-config');
|
|
if (!positive(b.tracker?.sync?.botId)) throw new BusError('tracker-config');
|
|
const roles = {};
|
|
for (const [role, r] of Object.entries(b.roles ?? {}))
|
|
if (r.tracker) {
|
|
if (!positive(r.tracker.botId) || typeof r.definition !== 'string') throw new BusError('tracker-config');
|
|
roles[role] = { definition: r.definition, botId: r.tracker.botId };
|
|
}
|
|
const pms = Object.keys(roles).filter((r) => roles[r].definition === 'pm');
|
|
if (pms.length !== 1) throw new BusError('tracker-config');
|
|
const labels = new Map();
|
|
for (const [title, id] of Object.entries(b.tracker.labels ?? {})) {
|
|
if (!positive(id)) throw new BusError('tracker-config');
|
|
labels.set(id, title);
|
|
}
|
|
return { business: name, baseUrl: t.baseUrl, project: t.project, pollSeconds, reconcileMinutes, roles, pm: pms[0], labels };
|
|
}
|
|
// Files behind the business's Vikunja references, to notice a token rotated on disk.
|
|
function files(b) {
|
|
const out = {};
|
|
for (const [role, r] of Object.entries(b.roles ?? {})) if (r.credentials?.vikunja?.file) out[role] = r.credentials.vikunja.file;
|
|
if (b.tracker?.sync?.credentials?.vikunja?.file) out['@sync'] = b.tracker.sync.credentials.vikunja.file;
|
|
return out;
|
|
}
|
|
const stamp = (file) => {
|
|
try {
|
|
const s = lstatSync(file);
|
|
return `${s.ino}:${s.size}:${s.mtimeMs}`;
|
|
} catch {
|
|
return 'missing';
|
|
}
|
|
};
|
|
export function tasksAdapter({
|
|
trackers = {},
|
|
fetch = globalThis.fetch,
|
|
clock = Date.now,
|
|
timers = globalThis,
|
|
log = (line) => process.stderr.write(line + '\n'),
|
|
autostart = true,
|
|
} = {}) {
|
|
return async function factory({ broker, credentials, businesses }) {
|
|
const contexts = new Map();
|
|
for (const [name, t] of Object.entries(trackers)) {
|
|
const s = settings(name, t, businesses[name]);
|
|
const call = client({ baseUrl: s.baseUrl, fetch });
|
|
const use = (instance) => (method, path, opts) =>
|
|
credentials.use(`${name}/${instance}`, 'vikunja', (token) => call(token, method, path, opts));
|
|
const ctx = {
|
|
...s,
|
|
broker,
|
|
credentials,
|
|
clock,
|
|
call,
|
|
per: 50,
|
|
view: null,
|
|
bots: new Set([businesses[name].tracker.sync.botId, ...Object.values(s.roles).map((r) => r.botId)]),
|
|
sync: use('@sync'),
|
|
as: (role, method, path, opts) => use(role)(method, path, opts),
|
|
state: 'starting',
|
|
refused: null,
|
|
queue: Promise.resolve(),
|
|
pending: 0,
|
|
lastStart: null,
|
|
lastReconcile: null,
|
|
lastComment: new Map(),
|
|
commentsAt: new Map(),
|
|
commentCount: new Map(),
|
|
files: Object.fromEntries(Object.entries(files(businesses[name])).map(([i, f]) => [i, { file: f, stamp: stamp(f) }])),
|
|
told: new Set(),
|
|
timers: [],
|
|
};
|
|
ctx.syncCall = (_token, method, path, opts) => ctx.sync(method, path, opts);
|
|
ctx.asCall = (role) => (_token, method, path, opts) => ctx.as(role, method, path, opts);
|
|
ctx.handle = verbs(ctx);
|
|
contexts.set(name, ctx);
|
|
}
|
|
let closed = false;
|
|
// One operation at a time per business: a tick never reads mid-write.
|
|
const serial = (ctx, fn) => {
|
|
const run = ctx.queue.then(() => (closed ? Promise.reject(new BusError('tracker-closed')) : fn()));
|
|
ctx.queue = run.catch(() => {});
|
|
return run;
|
|
};
|
|
const background = (ctx, what, fn) => {
|
|
if (ctx.pending > 1) return Promise.resolve();
|
|
ctx.pending++;
|
|
return serial(ctx, fn)
|
|
.catch((e) => {
|
|
if (!closed) log(`tasks ${ctx.business} ${what} ${code(e)}`);
|
|
})
|
|
.finally(() => ctx.pending--);
|
|
};
|
|
// credential.expiring, .expired and .changed, each once per instance per process.
|
|
function notify(ctx) {
|
|
const events = [];
|
|
const tell = (kind, instance, date) => {
|
|
const key = kind + ' ' + instance;
|
|
if (ctx.told.has(key)) return;
|
|
ctx.told.add(key);
|
|
events.push({ kind, body: { service: 'vikunja', instance, ...(date ? { date } : {}) } });
|
|
};
|
|
for (const s of credentials.status())
|
|
if (s.service === 'vikunja' && s.instance.startsWith(ctx.business + '/') && s.state !== 'valid')
|
|
tell(`credential.${s.state}`, s.instance, s.date);
|
|
for (const [instance, f] of Object.entries(ctx.files))
|
|
if (stamp(f.file) !== f.stamp) tell('credential.changed', `${ctx.business}/${instance}`, null);
|
|
if (events.length) broker.recordTask({ business: ctx.business, events });
|
|
}
|
|
async function start(ctx) {
|
|
try {
|
|
await startup(ctx);
|
|
await reconcile(ctx);
|
|
ctx.state = 'ready';
|
|
ctx.refused = null;
|
|
if (ctx.untested) log(`tasks ${ctx.business} startup untested-version`);
|
|
} catch (e) {
|
|
ctx.state = 'refused';
|
|
ctx.refused = code(e);
|
|
log(`tasks ${ctx.business} startup ${ctx.refused}`);
|
|
}
|
|
try {
|
|
notify(ctx);
|
|
} catch (e) {
|
|
log(`tasks ${ctx.business} credentials ${code(e)}`);
|
|
}
|
|
}
|
|
const api = {
|
|
timeout: TIMEOUT,
|
|
ready: null,
|
|
async handle(cap, verb, args) {
|
|
const { business } = broker.identity(cap);
|
|
const ctx = contexts.get(business);
|
|
if (!ctx) throw new BusError('tracker-unconfigured');
|
|
if (ctx.state !== 'ready') throw new BusError(ctx.state === 'starting' ? 'tracker-starting' : ctx.refused);
|
|
return serial(ctx, () => ctx.handle(cap, verb, args));
|
|
},
|
|
// For tests and the host's status verb. Each returns once the operation has run.
|
|
start: (business) => serial(contexts.get(business), () => start(contexts.get(business))),
|
|
tick: (business) => serial(contexts.get(business), () => tick(contexts.get(business))),
|
|
reconcile: (business) => serial(contexts.get(business), () => reconcile(contexts.get(business))),
|
|
notify: (business) => serial(contexts.get(business), async () => notify(contexts.get(business))),
|
|
status: () =>
|
|
[...contexts.values()].map((c) => ({
|
|
business: c.business,
|
|
state: c.state,
|
|
refused: c.refused,
|
|
version: c.version ?? null,
|
|
untested: c.untested ?? null,
|
|
lastPoll: c.lastStart === null ? null : new Date(c.lastStart).toISOString(),
|
|
lastReconcile: c.lastReconcile === null ? null : new Date(c.lastReconcile).toISOString(),
|
|
})),
|
|
async close() {
|
|
closed = true;
|
|
for (const c of contexts.values()) for (const t of c.timers) timers.clearInterval(t);
|
|
await Promise.all([...contexts.values()].map((c) => c.queue));
|
|
},
|
|
};
|
|
const every = (ctx, ms, fn) => {
|
|
const t = timers.setInterval(fn, ms);
|
|
t?.unref?.();
|
|
ctx.timers.push(t);
|
|
};
|
|
if (autostart) {
|
|
api.ready = Promise.all(
|
|
[...contexts.values()].map((ctx) => {
|
|
every(ctx, ctx.pollSeconds * 1000, () => {
|
|
if (ctx.state === 'ready')
|
|
background(ctx, 'tick', async () => {
|
|
await tick(ctx);
|
|
notify(ctx);
|
|
});
|
|
else if (ctx.state === 'refused' && TRANSIENT.has(ctx.refused)) background(ctx, 'startup', () => start(ctx));
|
|
});
|
|
every(ctx, ctx.reconcileMinutes * 60000, () => {
|
|
if (ctx.state === 'ready') background(ctx, 'reconcile', () => reconcile(ctx));
|
|
});
|
|
return background(ctx, 'startup', () => start(ctx));
|
|
}),
|
|
);
|
|
}
|
|
return api;
|
|
};
|
|
}
|