import { BusError } from '../../bus/src/broker.mjs'; import { outcome, refusal, all } from './vikunja.mjs'; import { ref, checkTask, taskFields, digest, tombstone } from './digest.mjs'; // The sync bot's reads (addendum B sections 3 and 5). A tick reads the whole open board, because a // move between open buckets doesn't bump `updated`, then the cursor: every task updated since the // previous tick's start minus WINDOW_MS. The window covers the async bump lag and the one-second // rounding (probes.md, U and Paging). Overlap is deduped by digest against task_current. export const WINDOW_MS = 60000; const iso = (ms) => new Date(ms).toISOString(); const floorSecond = (ms) => iso(Math.floor(ms / 1000) * 1000); // `updated` comes back to the second, or with nanoseconds on a create. task_external_changes compares // it as text, so every snapshot stores one form. export const normal = (updated) => iso(Date.parse(updated)); const shape = () => { throw new BusError('tracker-shape'); }; // The open board, every page. A bucket's tasks are null when it is empty; `count` is the bucket's total. export async function board(ctx) { const seen = new Map(); for (let page = 1; page <= 1000; page++) { const r = await ctx.sync('GET', `/projects/${ctx.project}/views/${ctx.view.id}/buckets/tasks`, { query: { filter: 'done = false', per_page: ctx.per, page }, }); if (outcome(r) !== 'ok') throw new BusError(refusal(r)); if (!Array.isArray(r.json?.items)) shape(); let more = false; for (const b of r.json.items) { const tasks = b?.tasks ?? []; if (!Number.isSafeInteger(b?.id) || !Array.isArray(tasks)) shape(); for (const t of tasks) { checkTask(t); seen.set(t.id, { task: t, bucket: b.id }); } if (Number.isSafeInteger(b.count) ? b.count > page * ctx.per : tasks.length >= ctx.per) more = true; } if (!more) return seen; } throw new BusError('tracker-paging'); } // A list walk keyed by id, so a write landing mid-walk can't shift a page under it. export async function walk(ctx, filter, query = {}) { const items = []; let last = 0; for (let i = 0; i < 100000; i++) { const r = await ctx.sync('GET', `/projects/${ctx.project}/tasks`, { query: { ...query, filter: [filter, `id > ${last}`].filter(Boolean).join(' && '), sort_by: 'id', order_by: 'asc', per_page: ctx.per, page: 1, }, }); if (outcome(r) !== 'ok') throw new BusError(refusal(r)); if (!Array.isArray(r.json?.items)) shape(); for (const t of r.json.items) { checkTask(t); if (t.id <= last || t.project_id !== ctx.project) shape(); last = t.id; items.push(t); } if (r.json.items.length < ctx.per) return items; } throw new BusError('tracker-paging'); } const changed = (before, after) => Object.keys(after).filter((k) => JSON.stringify(before?.[k]) !== JSON.stringify(after[k])); // Comments by anyone but this business's bots, newer than the last one seen. The first look at a // task counts only comments created inside the window, so a restart doesn't replay old ones. // `created` is to the second, so the bound is too. async function foreignComments(ctx, id, sinceMs) { const r = await all(ctx.syncCall, null, `/tasks/${id}/comments`, {}, ctx.per); if (!r.ok) throw new BusError(refusal(r.r)); const since = Math.floor(sinceMs / 1000) * 1000; const last = ctx.lastComment.get(id); let n = 0, top = last ?? 0; for (const c of r.items) { if (!Number.isSafeInteger(c?.id) || !Number.isSafeInteger(c.author?.id)) shape(); top = Math.max(top, c.id); const fresh = last === undefined ? Date.parse(c.created) >= since : c.id > last; if (fresh && !ctx.bots.has(c.author.id)) n++; } ctx.lastComment.set(id, top); return n; } // One task read. Records nothing that a newer self snapshot already covers, and nothing that // matches task_current; a comment count alone is an event without a snapshot. function consider(ctx, current, { task, bucket, via, readAt, comments = 0 }, out) { const r = ref(ctx.project, task.id); const fields = taskFields(task, bucket); const d = digest(fields); const c = current.get(r); const stale = c?.source === 'self' && readAt <= c.at; if (stale || c?.digest === d) { if (comments) out.events.push({ kind: 'task.changed.external', subject: r, body: { via, digest: c.digest, previous: c.digest, changed: [], comments }, }); return; } out.snapshots.push({ task_ref: r, updated: normal(task.updated), digest: d, fields, via, read_at: readAt }); out.events.push({ kind: 'task.changed.external', subject: r, body: { via, digest: d, previous: c?.digest ?? null, changed: changed(c?.fields, fields), ...(comments ? { comments } : {}) }, }); } // Open tasks the board no longer shows: done, moved, deleted or no longer shared. async function missing(ctx, current, refs, out) { for (const r of refs) { const c = current.get(r); const id = Number(r.slice(r.indexOf('/') + 1)); const readAt = iso(ctx.clock()); const x = await ctx.sync('GET', `/tasks/${id}`); let fields, updated, event; if (x.status === 200) { checkTask(x.json); updated = normal(x.json.updated); if (x.json.project_id !== ctx.project) { fields = tombstone('moved', x.json.project_id); event = { reason: 'moved', project: x.json.project_id }; } else if (x.json.done) fields = taskFields(x.json, ctx.view.done); else continue; // still open in the project: the board read raced a move, the next tick sees it } else if (outcome(x) === 'not-found' || x.status === 403) { updated = c.updated; fields = tombstone(x.status === 403 ? 'no-access' : 'not-found'); event = { reason: fields.gone }; } else throw new BusError(refusal(x)); const d = digest(fields); if (c.source === 'self' && readAt <= c.at) continue; out.snapshots.push({ task_ref: r, updated, digest: d, fields, via: 'task', read_at: readAt }); out.events.push( event ? { kind: 'task.missing', subject: r, body: event } : { kind: 'task.changed.external', subject: r, body: { via: 'task', digest: d, previous: c.digest, changed: changed(c.fields, fields) }, }, ); } } const record = (ctx, out) => { if (out.snapshots.length || out.events.length) ctx.broker.recordTask({ business: ctx.business, ...out }); }; export async function tick(ctx) { const startMs = ctx.clock(); const readAt = iso(startMs); const sinceMs = Math.floor(((ctx.lastStart ?? startMs) - WINDOW_MS) / 1000) * 1000; const open = await board(ctx); const hits = new Map((await walk(ctx, `updated > ${floorSecond(sinceMs)}`)).map((t) => [t.id, t])); const current = new Map(ctx.broker.taskView(ctx.business, 'current').map((x) => [x.task_ref, x])); const out = { snapshots: [], events: [] }; const comments = new Map(); for (const [id, t] of hits) { const key = normal(t.updated); if (ctx.commentsAt.get(id) === key) continue; comments.set(id, await foreignComments(ctx, id, sinceMs)); ctx.commentsAt.set(id, key); } for (const [id, { task, bucket }] of open) { const hit = hits.get(id); const use = hit && Date.parse(hit.updated) >= Date.parse(task.updated) ? hit : task; const b = use.done ? ctx.view.done : bucket; consider(ctx, current, { task: use, bucket: b, via: hit ? 'cursor' : 'board', readAt, comments: comments.get(id) }, out); } const explained = new Set(); for (const [id, t] of hits) { if (open.has(id)) continue; // An open task the board didn't show raced a move; the next board read places it. if (!t.done) continue; explained.add(ref(ctx.project, id)); consider(ctx, current, { task: t, bucket: ctx.view.done, via: 'cursor', readAt, comments: comments.get(id) }, out); } const gone = ctx.broker .taskView(ctx.business, 'open') .map((x) => x.task_ref) .filter((r) => !explained.has(r) && !open.has(Number(r.slice(r.indexOf('/') + 1)))); await missing(ctx, current, gone, out); record(ctx, out); ctx.lastStart = startMs; return { snapshots: out.snapshots.length, events: out.events.length }; } // The hourly full read, and the first read at startup. comment_count catches comments on tasks the // cursor window missed; a changed count is checked against the comment list before it is reported. export async function reconcile(ctx) { const startMs = ctx.clock(); const readAt = iso(startMs); const sinceMs = ctx.lastReconcile ?? startMs; const open = await board(ctx); const list = await walk(ctx, '', { expand: 'comment_count' }); const current = new Map(ctx.broker.taskView(ctx.business, 'current').map((x) => [x.task_ref, x])); const out = { snapshots: [], events: [] }; const listed = new Set(); for (const t of list) { listed.add(ref(ctx.project, t.id)); const bucket = t.done ? ctx.view.done : open.get(t.id)?.bucket; if (!bucket) continue; let comments = 0; const count = t.comment_count; if (count !== undefined && !Number.isSafeInteger(count)) shape(); const before = ctx.commentCount.get(t.id); if (count !== undefined) ctx.commentCount.set(t.id, count); if (before !== undefined && count !== before) comments = await foreignComments(ctx, t.id, sinceMs); consider(ctx, current, { task: t, bucket, via: 'reconcile', readAt, comments }, out); } const gone = [...current.values()] .filter((x) => x.fields.gone === undefined && !listed.has(x.task_ref)) .map((x) => x.task_ref); await missing(ctx, current, gone, out); record(ctx, out); ctx.lastReconcile = startMs; ctx.lastStart ??= startMs; return { snapshots: out.snapshots.length, events: out.events.length }; }