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: // {: {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; }; }