From 9294f01d646b8522a3284fe3b76398f40bf44140 Mon Sep 17 00:00:00 2001 From: Warren Date: Tue, 1 Sep 2026 04:47:26 +0000 Subject: [PATCH 1/2] wip: dispatch job + pure planner --- src/functions/index.ts | 25 ++- src/jobs/dispatch.job.ts | 454 ++++++++++++++++++++++++++++++++++++++ src/jobs/dispatch.plan.ts | 419 +++++++++++++++++++++++++++++++++++ src/jobs/index.ts | 14 +- 4 files changed, 906 insertions(+), 6 deletions(-) create mode 100644 src/jobs/dispatch.job.ts create mode 100644 src/jobs/dispatch.plan.ts diff --git a/src/functions/index.ts b/src/functions/index.ts index 32831e5..dd37a62 100644 --- a/src/functions/index.ts +++ b/src/functions/index.ts @@ -1,8 +1,23 @@ // Copyright (c) 2026 ObjectStack. Licensed under the Apache-2.0 license. // -// Named callables a `script` flow node resolves by name at run time. A flow -// function is PURE: it takes `input` and RETURNS a value; a later declarative -// node persists it. Add yours to this map — the map itself is already wired -// into objectstack.config.ts, so no feature branch needs to touch the config. +// Named callables the platform resolves BY NAME at run time — a `script` flow +// node's function, and a `defineJob`'s `handler`. Both are read from the same +// place: `collectBundleFunctions(bundle)` over `defineStack({ functions })`. +// Add yours to this map — the map itself is already wired into +// objectstack.config.ts, so no feature branch needs to touch the config. +// +// A flow-node function is PURE: it takes `input` and RETURNS a value, and a +// later declarative node persists it. A function that legitimately writes +// DECLARES it (`{ handler, effect: 'writes' }`) so a run reports +// `unmeasuredEffect` rather than claiming it wrote nothing. + +import { DISPATCH_HANDLER_NAME, dulyDispatch } from '../jobs/dispatch.job.js'; -export const dulyFunctions = {}; +export const dulyFunctions = { + // The dispatch job's handler. Declared `writes` because it does: it inserts + // `duly_task` rows directly rather than returning a value for a declarative + // node to persist. That is the honest declaration for a JOB handler, which + // has no flow graph around it to do the writing — see the note on data reach + // in `src/jobs/dispatch.job.ts`. + [DISPATCH_HANDLER_NAME]: { handler: dulyDispatch, effect: 'writes' as const }, +}; diff --git a/src/jobs/dispatch.job.ts b/src/jobs/dispatch.job.ts new file mode 100644 index 0000000..be4602a --- /dev/null +++ b/src/jobs/dispatch.job.ts @@ -0,0 +1,454 @@ +// Copyright (c) 2026 ObjectStack. Licensed under the Apache-2.0 license. + +import { defineJob } from '@objectstack/spec'; + +import { + DISPATCH_DUTY_FIELDS, + DISPATCHABLE_FORM, + DISPATCHABLE_STATUS, + FAULT_SKIP_REASONS, + nextDispatchedPeriod, + planDispatch, + type BackfillWindow, + type DispatchDuty, + type DutySkip, + type TaskDraft, +} from './dispatch.plan.js'; + +/** + * `duly_dispatch` — the spine. Active recurring duties become tasks, once per + * period, forever, without a lock. + * + * This file is the SMALL half. Every decision the dispatcher makes lives in + * `dispatch.plan.ts` as a pure function; what is left here is the schedule + * declaration and the code that talks to the engine. That boundary is the point + * — the imperative surface of the spine is about forty lines, and they are all + * on this page. + * + * ── What is declarative, and what could not be ─────────────────────────── + * Declarative: the schedule, the timezone, the retry policy, the timeout, the + * duty selection filter, and the identity constraint that makes the whole thing + * idempotent (`duly_task_dispatch_identity`, declared on the object). Imperative: + * the period arithmetic — calendar maths in a per-row IANA zone, which no + * filter language expresses — and the insert loop. + * + * The insert loop is imperative for a reason worth writing down, because the + * declarative alternative exists and was measured. A scheduled FLOW could do + * this shape (`get_record` → `script` → `loop` → `create_record`), and its + * error handling — a `try_catch` region or a `fault` edge — would swallow the + * duplicate insert. But it swallows every OTHER failure identically: the + * `create_record` executor collapses the engine's error to a string + * (`create_record(duly_task) failed: `) with no code, so a missing + * required field, a permission refusal and a store outage all arrive as the + * same anonymous failure the catch region cannot tell from a duplicate. For the + * one job the product cannot afford to fail quietly, "the run succeeded and + * created nothing" is the wrong failure to make easy. Two further measured + * facts settled it: the schedule trigger hands a flow + * `{ event, params: { jobId, flowName, schedule } }` and no input channel, so + * backfill's `{ from, to }` cannot reach a scheduled flow at all, while + * `IJobService.trigger(name, data)` forwards `data` to a job handler; and + * `ScheduleTrigger` calls `jobService.schedule(name, schedule, handler)` with no + * options, so a scheduled flow gets neither `retryPolicy` nor `timeout`. + * + * ── The one thing the platform could not give this job: data reach ─────── + * Measured against @objectstack/runtime 17.2.0 and @objectstack/service-job + * 17.2.0, a declarative job handler is invoked with exactly + * `{ jobId: string, data: unknown, bundle: object }`. There is no engine, no + * service registry and no logger on it, and `bundle` is the metadata bundle. + * So the one metadata shape the platform offers for scheduled work cannot, by + * itself, read or write a record. + * + * That is filed as a platform defect, not worked around silently. What this + * file does instead is make the seam explicit and singular: {@link runDispatch} + * takes its engine as an ARGUMENT — no module-scope client reached for at a + * distance, no `ctx.ql ?? ctx.engine ?? …` chain probing for a key the contract + * does not declare — and {@link bindDispatchEngine} is the one place a host + * supplies it. Until a host calls it, {@link dulyDispatch} REFUSES loudly with + * a message naming the wiring. A dispatcher that silently does nothing is the + * single worst failure this product can have; a job run that fails with "the + * dispatcher has no engine" is a bad night, not a silent quarter. + */ + +/** The job's metadata name. Exported so the metadata and the tests agree. */ +export const DISPATCH_JOB_NAME = 'duly_dispatch'; + +/** + * The `defineStack({ functions })` key the job's `handler` resolves against. + * + * `AppPlugin` looks the handler up as `collectBundleFunctions(bundle)[job.handler]` + * and, on a miss, logs `job handler not found in bundle.functions — skipping` + * and moves on. The job then reads as wired everywhere: registered in the + * metadata registry, listed in the admin UI, and never once executed. There is + * no author-time gate for it, which is why the name is a constant and + * `test/dispatch.test.ts` performs the same lookup the runtime performs. + */ +export const DISPATCH_HANDLER_NAME = 'dulyDispatch'; + +/** + * Duties read per round trip. + * + * The sweep pages rather than taking one large `limit`: an unbounded read is + * the thing that works on every test fixture and falls over on the first real + * tenant, and it fails by silently dispatching a PREFIX of the duties — the + * people whose ids sort late simply stop getting tasks, with no error anywhere. + */ +export const DISPATCH_PAGE_SIZE = 200; + +// ───────────────────────────────────────────────────────────────────────── +// The engine seam +// ───────────────────────────────────────────────────────────────────────── + +/** + * The slice of `IDataEngine` the dispatcher uses — structural, so the real + * engine satisfies it without an import and a test fake satisfies it without a + * kernel. + */ +export interface DispatchEngine { + find(objectName: string, query?: Record, options?: Record): Promise; + insert(objectName: string, data: Record, options?: Record): Promise; + update(objectName: string, data: Record, options?: Record): Promise; +} + +let boundEngine: DispatchEngine | null = null; + +/** + * Give the dispatch job its data engine. + * + * A host calls this once, from `defineStack({ onEnable })`, which is the only + * place in an ObjectStack application that is handed `ctx.ql`. It exists + * because the job handler context does not carry an engine (see the file + * header); when the platform grows one, this function and its caller are what + * gets deleted, and nothing else in the dispatcher changes. + */ +export function bindDispatchEngine(engine: DispatchEngine): void { + boundEngine = engine; +} + +/** Drop the binding. For tests, so one suite cannot leak an engine into the next. */ +export function unbindDispatchEngine(): void { + boundEngine = null; +} + +function requireDispatchEngine(): DispatchEngine { + if (boundEngine === null) { + throw new Error( + `Job '${DISPATCH_JOB_NAME}' has no data engine. The platform invokes a job handler with ` + + `{ jobId, data, bundle } and no engine handle, so the host must call ` + + `bindDispatchEngine(ctx.ql) from defineStack({ onEnable }). No tasks were dispatched.`, + ); + } + return boundEngine; +} + +// ───────────────────────────────────────────────────────────────────────── +// Backfill input +// ───────────────────────────────────────────────────────────────────────── + +const ISO_DATE = /^\d{4}-\d{2}-\d{2}$/; + +/** + * Read the optional `{ from?, to? }` job input into a backfill window. + * + * `IJobService.trigger(name, data)` forwards `data` to the handler unchanged, + * so this is the operator's channel — and the only one that is not a schedule. + * + * Three readings, all decided rather than inferred: + * + * - **Neither** → `null`, the scheduled run. This is what the cron produces. + * - **`from` alone** → `from` … the run's clock. "Catch this duty up." + * - **`to` alone** → REFUSED. It could only mean "from the beginning of the + * duty", and `effective_from` is nullable, so the honest reading of the + * unbounded case is "generate tasks back to 1583". A job that quietly + * generates ten thousand rows because an operator left a field blank is not + * a job anyone runs twice; a refusal costs one retyped command. + * + * Malformed input throws. A backfill is typed by a human at the moment they run + * it, so a rejected run they can see beats a window silently reinterpreted. + */ +export function parseBackfillWindow(data: unknown, now: Date): BackfillWindow | null { + if (data == null) return null; + if (typeof data !== 'object' || Array.isArray(data)) { + throw new TypeError(`${DISPATCH_JOB_NAME}: job input must be an object of the shape { from?, to? }`); + } + + const raw = data as { from?: unknown; to?: unknown }; + if (raw.from === undefined && raw.to === undefined) return null; + + if (raw.from === undefined) { + throw new RangeError( + `${DISPATCH_JOB_NAME}: a backfill needs 'from'. A 'to' on its own would mean "every period since the ` + + `duty began", and effective_from is nullable, so that window has no floor.`, + ); + } + + const from = raw.from; + if (typeof from !== 'string' || !ISO_DATE.test(from)) { + throw new RangeError(`${DISPATCH_JOB_NAME}: 'from' must be a YYYY-MM-DD date, received ${JSON.stringify(from)}`); + } + + let to: string; + if (raw.to === undefined) { + to = now.toISOString().slice(0, 10); + } else if (typeof raw.to === 'string' && ISO_DATE.test(raw.to)) { + to = raw.to; + } else { + throw new RangeError(`${DISPATCH_JOB_NAME}: 'to' must be a YYYY-MM-DD date, received ${JSON.stringify(raw.to)}`); + } + + if (to < from) { + throw new RangeError(`${DISPATCH_JOB_NAME}: backfill window ends before it starts (${from} … ${to})`); + } + return { from, to }; +} + +// ───────────────────────────────────────────────────────────────────────── +// The run +// ───────────────────────────────────────────────────────────────────────── + +export interface DispatchOptions { + /** The run's clock. Defaults to now; supplied by tests and by a replay. */ + now?: Date; + /** A backfill window, or `null`/absent for the scheduled sweep. */ + window?: BackfillWindow | null; +} + +export interface DispatchResult { + duties: number; + /** Task rows this run inserted. */ + created: number; + /** Task rows that already existed — the normal steady state, not a problem. */ + existing: number; + /** Duties whose `last_dispatched_period` this run moved forward. */ + advanced: number; + skipped: DutySkip[]; + /** Set when the run finished without doing all of its work. */ + degradedReason?: string; +} + +/** + * The execution context for every read and write the job makes. + * + * `isSystem` because the dispatcher acts for the organisation, not for a + * person: `duly_duty` is `sharingModel: 'private'`, so a user-scoped sweep + * would see only the duties of whoever happened to trigger it and would + * silently dispatch a fraction of the tenant. + */ +const SYSTEM_CONTEXT = { isSystem: true } as const; + +/** + * Run one dispatch pass. + * + * The engine is a parameter and not a module-scope client, so this function is + * exactly as testable as the planner and there is no hidden global to leak + * between runs. + */ +export async function runDispatch(engine: DispatchEngine, options: DispatchOptions = {}): Promise { + const now = options.now ?? new Date(); + const window = options.window ?? null; + + const duties = await readDispatchableDuties(engine); + const plan = planDispatch({ duties, now, window }); + + const createdByDuty = new Map(); + let created = 0; + let existing = 0; + + for (const draft of plan.drafts) { + const outcome = await insertOnce(engine, draft); + if (outcome === 'created') { + created += 1; + const keys = createdByDuty.get(draft.duty); + if (keys) keys.push(draft.period_key); + else createdByDuty.set(draft.duty, [draft.period_key]); + } else { + existing += 1; + } + } + + const byId = new Map(duties.map((duty) => [duty.id, duty] as const)); + let advanced = 0; + for (const [dutyId, keys] of createdByDuty) { + const next = nextDispatchedPeriod(byId.get(dutyId)?.last_dispatched_period, keys); + if (next === null) continue; + // Writes `duly_duty`, never `duly_task`. `last_update_at` is the task's + // stagnation signal and this job must never touch it — a dispatcher that + // reset it would make every stalled task look freshly worked, and the + // signal going quiet raises no error anywhere. + await engine.update( + 'duly_duty', + { last_dispatched_period: next }, + { where: { id: dutyId }, multi: false, context: SYSTEM_CONTEXT }, + ); + advanced += 1; + } + + const faults = plan.skipped.filter((skip) => FAULT_SKIP_REASONS.includes(skip.reason)); + return { + duties: duties.length, + created, + existing, + advanced, + skipped: plan.skipped, + ...(faults.length > 0 + ? { + degradedReason: `${faults.length} duty/duties could not be dispatched: ${faults + .map((f) => `${f.duty} (${f.reason}${f.detail ? `: ${f.detail}` : ''})`) + .join('; ')}`, + } + : {}), + }; +} + +/** + * Every active recurring duty, paged. + * + * The filter is deliberately only the two zone-independent facts. The + * effective-window test is NOT pushed into the query: `effective_from` and + * `effective_to` are calendar days in the DUTY's own zone, and "today" is a + * different day in Auckland and in Los Angeles at the same instant, so a single + * SQL predicate cannot be right for every row it matches. The window test + * belongs where the zone is known, which is the planner. + */ +async function readDispatchableDuties(engine: DispatchEngine): Promise { + const duties: DispatchDuty[] = []; + for (let offset = 0; ; offset += DISPATCH_PAGE_SIZE) { + const page = (await engine.find( + 'duly_duty', + { + where: { status: DISPATCHABLE_STATUS, form: DISPATCHABLE_FORM }, + fields: [...DISPATCH_DUTY_FIELDS], + // A stable total order, or paging re-reads and skips rows. + orderBy: [{ field: 'id', order: 'asc' }], + limit: DISPATCH_PAGE_SIZE, + offset, + }, + { context: SYSTEM_CONTEXT }, + )) as DispatchDuty[]; + duties.push(...page); + if (page.length < DISPATCH_PAGE_SIZE) return duties; + } +} + +/** + * Insert one task, or establish that it is already there. + * + * ── Attempt the insert; do not read first ──────────────────────────────── + * The happy path is a bare `insert`. A read-then-write guard would be a race + * two overlapping runs lose — both read nothing, both write, and the loser + * either duplicates the obligation or fails anyway — and it would cost a query + * on every task on every night for a collision that is rare. + * + * `duly_task_dispatch_identity`, unique on `(duty, owner, period_key)` and + * scoped to the organization, is what makes the bare insert safe. It is the + * constraint that lets this job be ordinary: re-runnable, killable mid-run, + * safe to invoke twice concurrently. + * + * ── Why the FAILURE path reads instead of classifying the error ────────── + * The obvious shape — catch, ask "was that a uniqueness violation?", swallow if + * so — is not available honestly. Measured on 17.2.0, ObjectQL does not wrap a + * constraint violation in a platform error: the raw driver error propagates, so + * on SQLite the app sees `code: 'SQLITE_CONSTRAINT_UNIQUE'`, on Postgres it + * would see `23505`, on MySQL `ER_DUP_ENTRY`. The platform's own + * dialect-independent predicate for this, `isUniqueViolationError`, lives in + * `@objectstack/types`, which neither `@objectstack/spec` nor + * `@objectstack/runtime` re-exports and which an application therefore cannot + * reach. The remaining options were to hard-code one driver's spelling or to + * pattern-match an error message — a consumer growing tolerance for a producer + * that does not answer, which is how the wrong driver ships silently. + * + * So the failure path asks the DATA, which every driver answers the same way: + * *is the row there now?* If it is, the obligation exists exactly once and this + * run's job on it is done, whoever won the race. If it is not, the insert failed + * for a reason that is not "already dispatched" — a missing required field, a + * refused write, a store that is down — and the error is re-thrown so the run + * FAILS. That is the half a blanket swallow gets wrong. + * + * This is not read-then-write: nothing is read before the insert, so the happy + * path is one round trip and there is no window between a check and a write for + * a second run to slip through. + */ +async function insertOnce(engine: DispatchEngine, draft: TaskDraft): Promise<'created' | 'existing'> { + try { + await engine.insert('duly_task', { ...draft }, { context: SYSTEM_CONTEXT }); + return 'created'; + } catch (error) { + const rows = await engine.find( + 'duly_task', + { + where: { duty: draft.duty, owner: draft.owner, period_key: draft.period_key }, + fields: ['id'], + limit: 1, + }, + { context: SYSTEM_CONTEXT }, + ); + if (rows.length > 0) return 'existing'; + throw error; + } +} + +// ───────────────────────────────────────────────────────────────────────── +// The handler +// ───────────────────────────────────────────────────────────────────────── + +/** What the platform hands a job handler. Measured, not assumed — see the header. */ +export interface DispatchJobContext { + jobId?: string; + data?: unknown; +} + +/** + * What the run reports. + * + * A run that inserted nothing because everything already existed resolves + * `completed`: that is a SUCCESSFUL run, and the whole design rests on it being + * unremarkable. `degraded` is reserved for a run that finished without doing all + * of its work — a duty the period engine refused. It is not a failure and never + * retries; only a thrown error does that, which is what an unreachable store or + * an unbound engine produces. + */ +export const dulyDispatch = async ( + context: DispatchJobContext = {}, +): Promise<{ outcome: 'completed' | 'degraded'; reason?: string }> => { + const engine = requireDispatchEngine(); + const now = new Date(); + const window = parseBackfillWindow(context.data, now); + const result = await runDispatch(engine, { now, window }); + return result.degradedReason ? { outcome: 'degraded', reason: result.degradedReason } : { outcome: 'completed' }; +}; + +// ───────────────────────────────────────────────────────────────────────── +// The metadata +// ───────────────────────────────────────────────────────────────────────── + +/** + * Runs at 01:00 UTC daily, and resolves each duty in ITS OWN `timezone`. + * + * One UTC pass covers every zone because the question asked per duty is "is + * this period's task due to exist yet", never "is it midnight here". A job that + * chased local midnights would need one schedule per zone and would still get + * the answer wrong twice a year. + * + * 01:00 rather than 00:00 so a run is never racing a zone's own DST transition, + * and so the hour every other nightly job in every other system picks is left + * alone. + */ +export const DispatchJob = defineJob({ + name: DISPATCH_JOB_NAME, + label: 'Dispatch due tasks', + description: + 'Turns active recurring duties into tasks, one per period, in each duty\'s own timezone. Idempotent on (duty, owner, period_key); safe to re-run, backfill or interrupt.', + schedule: { type: 'cron', expression: '0 1 * * *', timezone: 'UTC' }, + handler: DISPATCH_HANDLER_NAME, + + // A failed run is worth retrying: the usual cause is a store that was briefly + // unreachable, and the identity index makes a retry after a partial run cost + // nothing — the tasks already written are found, not duplicated. Three + // attempts over a few minutes, then leave it for the next night rather than + // hammering a store that is genuinely down. + retryPolicy: { maxRetries: 3, backoffMs: 30_000, backoffMultiplier: 2, maxRetryDelayMs: 300_000, jitter: true }, + + // Fifteen minutes. A dispatch pass is a paged read plus one insert per owed + // task; an attempt still running after fifteen minutes has stopped making + // progress, and abandoning it is safe for the same reason a retry is. + timeout: 900_000, +}); diff --git a/src/jobs/dispatch.plan.ts b/src/jobs/dispatch.plan.ts new file mode 100644 index 0000000..8e029da --- /dev/null +++ b/src/jobs/dispatch.plan.ts @@ -0,0 +1,419 @@ +// Copyright (c) 2026 ObjectStack. Licensed under the Apache-2.0 license. + +import { + FREQUENCIES, + dueDateFor, + periodBounds, + periodKeyFor, + periodsBetween, + visibleFromFor, + type DueAnchor, + type Frequency, +} from '../functions/period.js'; + +/** + * The dispatcher's BRAIN — pure, clockless, engineless. + * + * Everything the dispatcher decides lives here: which duties are in scope, + * which periods they owe a task for, and what each task row should say. It + * takes rows in and returns rows out. No I/O, no `Date.now()`, no platform + * import. `dispatch.job.ts` next door is the other half: the schedule + * declaration and the twenty lines that actually talk to the engine. + * + * The split is a file boundary on purpose. This module is the part that can be + * wrong in ways nothing would notice — a period key off by one zone, a lead + * window that never opens — and a pure function is the only shape that can be + * cornered by a test without booting anything. + * + * ── Every period spelling comes from `../functions/period.ts` ──────────── + * Not one key, boundary or due date is derived here. `duly_task` is unique on + * `(duty, owner, period_key)`, so two spellings of one period are two tasks for + * one obligation and nothing downstream can tell they were meant to be the + * same. Two consequences worth stating, because they read as coincidences: + * + * - **"today, in the duty's zone" is `periodKeyFor('daily', now, tz)`.** The + * daily period key IS the local calendar day, by that module's own spelling + * (`YYYY-MM-DD`). Asking it is how this file gets a local date without + * owning a second copy of the zone arithmetic. + * - **"the instant local day D begins in zone Z" is + * `periodBounds('daily', D, tz).start`.** Same reason. `periodsBetween` + * wants instants and a backfill window is authored as calendar dates; this + * is the conversion, performed by the module that owns it. + */ + +// ───────────────────────────────────────────────────────────────────────── +// Cadence defaults +// ───────────────────────────────────────────────────────────────────────── + +/** + * The cadence values used when a duty row carries `null`. + * + * These are NOT this module's opinion — each one is `duly_duty`'s own declared + * default, restated here so the planner stays a pure function with no metadata + * import, and pinned against the object schema in `test/dispatch.test.ts` so + * the two cannot drift into two answers. (Same pattern, and the same reason, as + * `DEFAULT_DUTY_TIMEZONE` in `src/actions/catalog.handlers.ts`.) + * + * The engine applies these on write, so a row read back normally carries them. + * The fallback covers the rows that predate a default or arrived through a path + * that skipped it — and it is a fallback to the DECLARED value, never an + * invented one. + */ +export const DEFAULT_TIMEZONE = 'UTC'; +export const DEFAULT_DUE_ANCHOR: DueAnchor = 'period_start'; +export const DEFAULT_DUE_OFFSET_DAYS = 0; +export const DEFAULT_LEAD_DAYS = 7; + +/** The status a duty must hold to dispatch. */ +export const DISPATCHABLE_STATUS = 'active'; +/** The form a duty must hold to dispatch. `standing` and `one_off` never do. */ +export const DISPATCHABLE_FORM = 'recurring'; + +/** + * The `duly_duty` projection the planner reads. + * + * Exported so the read in `dispatch.job.ts` and the fields consumed here are + * one list rather than two. A field the planner reads but the projection omits + * comes back `undefined`, which reads exactly like "the duty does not set it" — + * a silent wrong answer, not an error. `test/dispatch.test.ts` pins the two + * together. + */ +export const DISPATCH_DUTY_FIELDS = [ + 'id', + 'name', + 'form', + 'status', + 'owner', + 'business_unit', + 'source', + 'frequency', + 'due_anchor', + 'due_offset_days', + 'lead_days', + 'timezone', + 'effective_from', + 'effective_to', + 'last_dispatched_period', +] as const; + +// ───────────────────────────────────────────────────────────────────────── +// Shapes +// ───────────────────────────────────────────────────────────────────────── + +/** A `duly_duty` row, as the projection above returns it. */ +export interface DispatchDuty { + id: string; + name?: string | null; + form?: string | null; + status?: string | null; + owner?: string | null; + business_unit?: string | null; + source?: string | null; + frequency?: string | null; + due_anchor?: string | null; + due_offset_days?: number | null; + lead_days?: number | null; + timezone?: string | null; + effective_from?: string | null; + effective_to?: string | null; + last_dispatched_period?: string | null; +} + +/** One `duly_task` row the dispatcher intends to insert. */ +export interface TaskDraft { + /** Copied at dispatch, so renaming the duty never rewrites history. */ + subject: string; + duty: string; + owner: string; + /** Denormalised from the DUTY, so a later transfer does not move rollups. */ + business_unit: string | null; + source: string; + period_key: string; + due_date: string; + visible_from: string; + status: 'open'; +} + +/** Why a duty produced no drafts on this run. */ +export type DispatchSkipReason = + /** `form: 'standing'` — a standing duty is attestable, never tickable. */ + | 'standing' + /** `form: 'one_off'` — dispatched by hand or by the assignment fan-out. */ + | 'one_off' + /** A form this dispatcher does not know. */ + | 'unknown_form' + /** `paused` or `retired`. */ + | 'not_active' + /** Recurring with no cadence. `recurring_needs_frequency` should prevent it. */ + | 'no_frequency' + /** A `frequency` value outside `FREQUENCIES`. */ + | 'unknown_frequency' + /** Today, in the duty's own zone, is outside `[effective_from, effective_to]`. */ + | 'outside_effective_window' + /** In scope, but no period is due to exist yet. */ + | 'nothing_due' + /** The period engine refused this duty's cadence — a bad zone, a bad offset. */ + | 'invalid_cadence'; + +/** + * The skip reasons that mean something is WRONG with the data, as opposed to + * the duty simply having nothing to do. A run that hits one of these finished, + * but did not do all of its work — which is what `JobRunOutcome.degraded` is + * for. The others are ordinary, expected, and silent. + */ +export const FAULT_SKIP_REASONS: readonly DispatchSkipReason[] = [ + 'unknown_form', + 'no_frequency', + 'unknown_frequency', + 'invalid_cadence', +]; + +export interface DutySkip { + duty: string; + reason: DispatchSkipReason; + /** The period engine's own message, for `invalid_cadence` only. */ + detail?: string; +} + +export interface DispatchPlan { + drafts: TaskDraft[]; + skipped: DutySkip[]; +} + +/** An inclusive backfill window, as authored calendar dates. */ +export interface BackfillWindow { + from: string; + to: string; +} + +export interface DispatchPlanInput { + duties: readonly DispatchDuty[]; + /** The run's clock. Supplied, never read from the environment. */ + now: Date; + /** Present for a backfill run; `null` for the scheduled run. */ + window?: BackfillWindow | null; +} + +// ───────────────────────────────────────────────────────────────────────── +// Helpers over the period engine +// ───────────────────────────────────────────────────────────────────────── + +/** + * The local calendar day containing `instant`, in `timezone`, as `YYYY-MM-DD`. + * + * This is the daily period key, which the period engine defines to be exactly + * the local calendar day. Asking for it that way is what keeps this file from + * owning a second copy of the zone arithmetic. + */ +function localDayOf(instant: Date, timezone: string): string { + return periodKeyFor('daily', instant, timezone); +} + +/** + * The first instant of local day `isoDate` in `timezone`. + * + * A daily period's bounds ARE the local day's bounds, so this is the period + * engine answering, not a conversion invented here. On a local day a zone + * erased at the date line the period holds no instants and `start === end`; + * `periodsBetween` walks past such a day, so feeding it one is safe. + */ +function startOfLocalDay(isoDate: string, timezone: string): Date { + return periodBounds('daily', isoDate, timezone).start; +} + +/** + * `isoDate` shifted forward by `days` calendar days. + * + * `visibleFromFor` is the period engine's civil-date shift — it subtracts its + * second argument — so a negative lead is a forward shift. Spelled through it + * rather than reimplemented because a second `addDays` is a second thing that + * can be wrong about February. + */ +function addCalendarDays(isoDate: string, days: number): string { + return visibleFromFor(isoDate, -days); +} + +function isFrequency(value: unknown): value is Frequency { + return typeof value === 'string' && (FREQUENCIES as readonly string[]).includes(value); +} + +/** `YYYY-MM-DD` values compare lexicographically, so a plain `<=` is a date test. */ +function withinDateWindow(day: string, from?: string | null, to?: string | null): boolean { + if (typeof from === 'string' && from !== '' && day < from) return false; + if (typeof to === 'string' && to !== '' && day > to) return false; + return true; +} + +// ───────────────────────────────────────────────────────────────────────── +// The plan +// ───────────────────────────────────────────────────────────────────────── + +/** + * Turn a page of duties into the task rows they owe. + * + * ── Which periods a SCHEDULED run considers ────────────────────────────── + * The current period, plus any later period whose lead window has already + * opened (`visible_from <= today`). Both halves are needed: + * + * - The current period unconditionally. `visible_from` governs when a task + * SHOWS UP, not whether it EXISTS — a monthly duty due on the 30th with no + * lead time still owes a row on the 1st, or nothing can be worked early. + * - Later periods, gated on visibility. A duty due on the 3rd with 7 days of + * lead is meant to appear in the previous month; if the dispatcher only ever + * emitted the current period, that row would not exist until the month it is + * due and the lead time would be decorative. + * + * The look-ahead is bounded at `today + lead_days`, which is exact rather than + * generous: `visible_from` is `due_date - lead_days` and `due_date` is clamped + * inside its own period, so a period starting after `today + lead_days` cannot + * have `visible_from <= today`. Nothing is missed and nothing extra is walked. + * + * ── Which periods a BACKFILL run considers ─────────────────────────────── + * Exactly the window, via `periodsBetween` — both ends inclusive of the period + * containing them. Visibility does not gate a backfill: an operator asking for + * March is asking for March, and every one of those periods is in the past, so + * its lead window opened long ago. + * + * ── The two effective-window tests, and why there are two ──────────────── + * 1. **Duty-level, on `today`** — the scheduled run's selection rule: a duty + * whose window has not opened or has closed is not dispatched at all. This + * is the rule as specified, and it is deliberately NOT applied to a + * backfill: backfilling a duty whose window closed last month is the whole + * point of backfill. + * 2. **Period-level, on `due_date`** — always. A period whose due date falls + * outside the duty's effective window is not owed, in either mode. This is + * what keeps a backfill from inventing obligations that predate the duty. + * + * ── One bad duty does not stop the run ─────────────────────────────────── + * A typo'd IANA zone or a non-integer offset makes the period engine throw. It + * throws per duty, so it is caught per duty and recorded as `invalid_cadence` + * with the engine's own message. The alternative — one malformed row stopping + * every other person's tasks from being created — is the worse failure by a + * wide margin. It is not swallowed: `FAULT_SKIP_REASONS` carries it into the + * job's `degraded` outcome, so the run is recorded as one that did not do all + * of its work. + */ +export function planDispatch(input: DispatchPlanInput): DispatchPlan { + const { duties, now, window = null } = input; + const drafts: TaskDraft[] = []; + const skipped: DutySkip[] = []; + + for (const duty of duties) { + const outcome = planForDuty(duty, now, window); + if ('reason' in outcome) { + skipped.push({ duty: duty.id, reason: outcome.reason, ...(outcome.detail ? { detail: outcome.detail } : {}) }); + continue; + } + drafts.push(...outcome.drafts); + } + + return { drafts, skipped }; +} + +type DutyOutcome = { drafts: TaskDraft[] } | { reason: DispatchSkipReason; detail?: string }; + +function planForDuty(duty: DispatchDuty, now: Date, window: BackfillWindow | null): DutyOutcome { + // ── Form and status: the two invariants, stated first ────────────────── + // A standing duty NEVER generates a task. Skipped entirely — never created + // and immediately closed, which would put an unclosable row in every list. + if (duty.form === 'standing') return { reason: 'standing' }; + if (duty.form === 'one_off') return { reason: 'one_off' }; + if (duty.form !== DISPATCHABLE_FORM) return { reason: 'unknown_form' }; + if (duty.status !== DISPATCHABLE_STATUS) return { reason: 'not_active' }; + + if (duty.frequency == null || duty.frequency === '') return { reason: 'no_frequency' }; + if (!isFrequency(duty.frequency)) return { reason: 'unknown_frequency' }; + const frequency: Frequency = duty.frequency; + + const timezone = duty.timezone ?? DEFAULT_TIMEZONE; + const dueAnchor = (duty.due_anchor ?? DEFAULT_DUE_ANCHOR) as DueAnchor; + const dueOffsetDays = duty.due_offset_days ?? DEFAULT_DUE_OFFSET_DAYS; + const leadDays = duty.lead_days ?? DEFAULT_LEAD_DAYS; + + try { + // "Today" is resolved in the DUTY's zone, not the server's. A single UTC + // pass covers every zone because the question asked per duty is "is this + // period's task due to exist yet", never "is it midnight here". + const today = localDayOf(now, timezone); + + if (window === null && !withinDateWindow(today, duty.effective_from, duty.effective_to)) { + return { reason: 'outside_effective_window' }; + } + + const currentKey = periodKeyFor(frequency, now, timezone); + const candidates = + window === null + ? periodsBetween(frequency, now, startOfLocalDay(addCalendarDays(today, leadDays), timezone), timezone) + : periodsBetween( + frequency, + startOfLocalDay(window.from, timezone), + startOfLocalDay(window.to, timezone), + timezone, + ); + + const drafts: TaskDraft[] = []; + for (const periodKey of candidates) { + const dueDate = dueDateFor({ frequency, periodKey, timezone, dueAnchor, dueOffsetDays }); + const visibleFrom = visibleFromFor(dueDate, leadDays); + + // Period-level effective clip — both modes. A period whose due date sits + // outside the duty's window is not owed by anyone. + if (!withinDateWindow(dueDate, duty.effective_from, duty.effective_to)) continue; + + // Visibility gate — scheduled runs only. The current period is exempt: + // it must exist whether or not it is meant to be on screen yet. + if (window === null && periodKey !== currentKey && visibleFrom > today) continue; + + drafts.push({ + subject: duty.name ?? '', + duty: duty.id, + owner: duty.owner ?? '', + business_unit: duty.business_unit ?? null, + source: duty.source ?? '', + period_key: periodKey, + due_date: dueDate, + visible_from: visibleFrom, + status: 'open', + }); + } + + if (drafts.length === 0) return { reason: 'nothing_due' }; + return { drafts }; + } catch (error) { + return { reason: 'invalid_cadence', detail: error instanceof Error ? error.message : String(error) }; + } +} + +// ───────────────────────────────────────────────────────────────────────── +// last_dispatched_period +// ───────────────────────────────────────────────────────────────────────── + +/** + * Where `duly_duty.last_dispatched_period` should stand after a run, or `null` + * when it should not be written at all. + * + * **Advances, never regresses.** The keys of one frequency are fixed-width and + * zero-padded (`2026-08` · `2026-W34` · `2026-Q3` · `2026-H2` · `2026` · + * `2026-08-21`), so a lexical `>` is chronological and no parsing is needed. + * A duty holds one frequency at a time, which is what makes that sound. + * + * Two consequences of the never-regress rule, stated so they are decisions: + * + * - A run that created nothing — every task already existed — writes nothing. + * That is the correct reading of "the latest key it created", and it is also + * what keeps a second identical run from issuing a single write. + * - If a duty's `frequency` is CHANGED, the stored key is in the old spelling + * and the comparison against the new one is meaningless. The rule then fails + * in the safe direction: the field either advances or is left alone, and it + * never jumps backwards. It is a bookkeeping breadcrumb, not the idempotency + * mechanism — that is the unique index — so a stale one costs nothing. + */ +export function nextDispatchedPeriod(current: string | null | undefined, createdKeys: readonly string[]): string | null { + let best: string | null = null; + for (const key of createdKeys) { + if (best === null || key > best) best = key; + } + if (best === null) return null; + if (typeof current === 'string' && current !== '' && current >= best) return null; + return best; +} diff --git a/src/jobs/index.ts b/src/jobs/index.ts index 61d701b..eb17c31 100644 --- a/src/jobs/index.ts +++ b/src/jobs/index.ts @@ -12,5 +12,17 @@ // resolves it against the keyed branch of `MetadataCollectionInput`, which // makes `name` optional and fails the assignment. A named array is `never[]` // while empty and infers correctly the moment something is pushed into it. +// +// ⚠ A job needs TWO registrations, and the second one has no author-time gate. +// This barrel publishes the SCHEDULE; `src/functions/index.ts` publishes the +// HANDLER the schedule names. `AppPlugin` resolves the handler as +// `collectBundleFunctions(bundle)[job.handler]` and, on a miss, logs +// "job handler not found in bundle.functions — skipping" and carries on: the +// job is registered, listed, and never executed. `test/dispatch.test.ts` +// performs the same lookup so the pair cannot come apart. + +import { DispatchJob } from './dispatch.job.js'; + +export { DispatchJob }; -export const dulyJobs = []; +export const dulyJobs = [DispatchJob]; From 20d1b2a2da2bbcfebd5531a24cacf1c07932d790 Mon Sep 17 00:00:00 2001 From: Warren Date: Tue, 1 Sep 2026 04:51:51 +0000 Subject: [PATCH 2/2] Add the dispatch job: idempotent task generation, backfill-capable --- src/jobs/dispatch.plan.ts | 12 +- test/dispatch.test.ts | 746 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 757 insertions(+), 1 deletion(-) create mode 100644 test/dispatch.test.ts diff --git a/src/jobs/dispatch.plan.ts b/src/jobs/dispatch.plan.ts index 8e029da..de31712 100644 --- a/src/jobs/dispatch.plan.ts +++ b/src/jobs/dispatch.plan.ts @@ -341,9 +341,19 @@ function planForDuty(duty: DispatchDuty, now: Date, window: BackfillWindow | nul } const currentKey = periodKeyFor(frequency, now, timezone); + // Both ends are the START of a local day, never `now`. `periodsBetween` + // returns `[]` when `to` is before `from`, and with `lead_days: 0` the + // horizon IS today — so passing the run's instant as `from` would put a + // mid-morning `now` after midnight-of-today and silently dispatch NOTHING, + // for every duty with no lead time, on every run. (Measured: it did.) const candidates = window === null - ? periodsBetween(frequency, now, startOfLocalDay(addCalendarDays(today, leadDays), timezone), timezone) + ? periodsBetween( + frequency, + startOfLocalDay(today, timezone), + startOfLocalDay(addCalendarDays(today, leadDays), timezone), + timezone, + ) : periodsBetween( frequency, startOfLocalDay(window.from, timezone), diff --git a/test/dispatch.test.ts b/test/dispatch.test.ts new file mode 100644 index 0000000..b2a7b6b --- /dev/null +++ b/test/dispatch.test.ts @@ -0,0 +1,746 @@ +// Copyright (c) 2026 ObjectStack. Licensed under the Apache-2.0 license. + +import { afterAll, afterEach, beforeAll, describe, expect, it } from 'vitest'; +import { AppPlugin, ObjectKernel, createStandaloneStack } from '@objectstack/runtime'; + +import stack from '../objectstack.config.js'; +import { Duty, Task } from '../src/objects/index.js'; +import { dulyJobs } from '../src/jobs/index.js'; +import { dulyFunctions } from '../src/functions/index.js'; +import { + DISPATCH_HANDLER_NAME, + DISPATCH_JOB_NAME, + DISPATCH_PAGE_SIZE, + DispatchJob, + bindDispatchEngine, + dulyDispatch, + parseBackfillWindow, + runDispatch, + unbindDispatchEngine, + type DispatchEngine, +} from '../src/jobs/dispatch.job.js'; +import { + DEFAULT_DUE_ANCHOR, + DEFAULT_DUE_OFFSET_DAYS, + DEFAULT_LEAD_DAYS, + DEFAULT_TIMEZONE, + DISPATCH_DUTY_FIELDS, + nextDispatchedPeriod, + planDispatch, + type DispatchDuty, +} from '../src/jobs/dispatch.plan.js'; +import { dueDateFor, periodKeyFor, periodsBetween, visibleFromFor } from '../src/functions/period.js'; + +/** + * The dispatch spine. + * + * Three layers, tested three ways, because each has a different way of being + * wrong invisibly: + * + * 1. **Wiring.** A job needs two registrations — the schedule in `dulyJobs` + * and the handler in `dulyFunctions` — and the second has no author-time + * gate at all. `pnpm validate` passes on a job whose handler name matches + * nothing; `AppPlugin` logs a warning at boot and skips it. So the lookup + * the runtime performs is performed here. + * + * 2. **The planner**, pure. Every scheduling decision, driven directly with + * fixed clocks and no engine. This is where the zone and lead-time + * behaviour is cornered, because a planner that is wrong about Los Angeles + * produces a task on the wrong day and nothing errors. + * + * 3. **Idempotency**, against a REAL booted engine on **sqlite**. Not the + * memory driver, and this is the load-bearing choice in the file: measured + * on 17.2.0, `InMemoryDriver.create` is a `table.push()` that stores no + * constraints of any kind, so two identical `duly_task` inserts BOTH + * SUCCEED and the table ends with two rows. A dispatcher suite on the + * memory driver would report idempotency passing while the index that + * provides it was never consulted. `test/task-hook.test.ts` uses memory + * legitimately — it tests hook ordering, which needs no constraint — but + * the whole point of this job is a constraint, so it has to run somewhere + * that has one. + */ + +type AnyRow = Record; + +// ───────────────────────────────────────────────────────────────────────── +// 1. Wiring — the failure mode that reads as success +// ───────────────────────────────────────────────────────────────────────── + +describe('wiring', () => { + it('is exported from the jobs barrel', () => { + expect(dulyJobs).toContain(DispatchJob); + expect(DispatchJob.name).toBe(DISPATCH_JOB_NAME); + }); + + it('reaches defineStack({ jobs }) — the only place the runtime reads', () => { + const names = ((stack as { jobs?: Array<{ name?: string }> }).jobs ?? []).map((j) => j.name); + expect(names).toContain(DISPATCH_JOB_NAME); + }); + + it("resolves its handler the way AppPlugin does: bundle.functions[job.handler]", () => { + // `collectBundleFunctions` reads `defineStack({ functions })` and takes + // `.handler` off each entry. A name that misses is logged and skipped — + // the job stays registered and never runs, with no author-time gate. + const functions = (stack as { functions?: Record }).functions ?? {}; + const entry = functions[DispatchJob.handler] as { handler?: unknown } | undefined; + expect(entry, `no function named '${DispatchJob.handler}' — the job would never execute`).toBeDefined(); + expect(typeof entry?.handler).toBe('function'); + expect(DispatchJob.handler).toBe(DISPATCH_HANDLER_NAME); + }); + + it('declares the handler as a writer, because it writes', () => { + // A `script`-node function is pure by contract; one that writes declares it + // so a run reports `unmeasuredEffect` instead of claiming it wrote nothing. + const entry = dulyFunctions[DISPATCH_HANDLER_NAME] as { effect?: string }; + expect(entry.effect).toBe('writes'); + }); + + it('is a UTC cron, and enabled', () => { + expect(DispatchJob.schedule.type).toBe('cron'); + const schedule = DispatchJob.schedule as { timezone?: string; expression?: unknown }; + expect(schedule.timezone).toBe('UTC'); + // A disabled job is registered and never scheduled, which reads identically + // to a working one from every surface except the logs. + expect(DispatchJob.enabled).not.toBe(false); + }); + + it('threads a retry policy and a timeout to the adapter', () => { + expect(DispatchJob.retryPolicy?.maxRetries).toBeGreaterThan(0); + expect(DispatchJob.timeout).toBeGreaterThan(0); + }); +}); + +describe('the cadence fallbacks are the object schema, not a second opinion', () => { + // The planner is pure and imports no metadata, so it restates `duly_duty`'s + // declared defaults. These assertions are what stop the two from drifting + // into two answers — the same pin as DEFAULT_DUTY_TIMEZONE in the catalog + // handlers. + const defaultOption = (field: { options?: Array<{ value: string; default?: boolean }> }) => + field.options?.find((o) => o.default)?.value; + + it('timezone', () => expect(Duty.fields.timezone.defaultValue).toBe(DEFAULT_TIMEZONE)); + it('lead_days', () => expect(Duty.fields.lead_days.defaultValue).toBe(DEFAULT_LEAD_DAYS)); + it('due_offset_days', () => expect(Duty.fields.due_offset_days.defaultValue).toBe(DEFAULT_DUE_OFFSET_DAYS)); + it('due_anchor', () => expect(defaultOption(Duty.fields.due_anchor)).toBe(DEFAULT_DUE_ANCHOR)); +}); + +describe('the duty projection covers every field the planner reads', () => { + it('names only real duly_duty fields', () => { + const declared = new Set([...Object.keys(Duty.fields), 'id']); + for (const field of DISPATCH_DUTY_FIELDS) { + expect(declared.has(field), `'${field}' is not a duly_duty field`).toBe(true); + } + }); + + it('includes every field a draft is built from', () => { + // A field the planner reads but the projection omits comes back undefined, + // which reads exactly like "the duty does not set it". + for (const field of [ + 'name', + 'form', + 'status', + 'owner', + 'business_unit', + 'source', + 'frequency', + 'due_anchor', + 'due_offset_days', + 'lead_days', + 'timezone', + 'effective_from', + 'effective_to', + 'last_dispatched_period', + ]) { + expect(DISPATCH_DUTY_FIELDS as readonly string[]).toContain(field); + } + }); +}); + +// ───────────────────────────────────────────────────────────────────────── +// 2. The planner — pure, fixed clocks +// ───────────────────────────────────────────────────────────────────────── + +const duty = (over: Partial = {}): DispatchDuty => ({ + id: 'duty_1', + name: 'File the emissions return', + form: 'recurring', + status: 'active', + owner: 'user_alice', + business_unit: 'bu_plant', + source: 'catalog', + frequency: 'monthly', + due_anchor: 'period_start', + due_offset_days: 4, + lead_days: 0, + timezone: 'UTC', + effective_from: null, + effective_to: null, + last_dispatched_period: null, + ...over, +}); + +const keysOf = (duties: DispatchDuty[], now: Date, window?: { from: string; to: string } | null) => + planDispatch({ duties, now, window: window ?? null }).drafts.map((d) => d.period_key); + +describe('what is never dispatched', () => { + const now = new Date('2026-08-15T09:00:00Z'); + + it('a standing duty produces nothing, in any run', () => { + // Not created and immediately closed — skipped entirely. A standing duty + // never completes, so a task for one is a row nobody can ever tick. + expect(keysOf([duty({ form: 'standing' })], now)).toEqual([]); + expect(keysOf([duty({ form: 'standing' })], now, { from: '2025-01-01', to: '2026-12-31' })).toEqual([]); + expect(planDispatch({ duties: [duty({ form: 'standing' })], now }).skipped[0]?.reason).toBe('standing'); + }); + + it('a one-off duty produces nothing — that is the fan-out\'s job', () => { + expect(keysOf([duty({ form: 'one_off' })], now)).toEqual([]); + expect(planDispatch({ duties: [duty({ form: 'one_off' })], now }).skipped[0]?.reason).toBe('one_off'); + }); + + it('a paused duty produces nothing, and a retired one produces nothing', () => { + for (const status of ['paused', 'retired']) { + expect(keysOf([duty({ status })], now), status).toEqual([]); + expect(planDispatch({ duties: [duty({ status })], now }).skipped[0]?.reason).toBe('not_active'); + } + }); + + it('un-pausing produces only the current period, never the gap', () => { + // The duty was last dispatched in February and has been paused since. The + // rule is "the current period", not "everything you missed" — a person + // coming back from three months' leave gets this month's work, not ninety + // days of backlog they cannot do anything about. + const resumed = duty({ status: 'active', last_dispatched_period: '2026-02' }); + expect(keysOf([resumed], now)).toEqual(['2026-08']); + }); + + it('a duty outside its effective window produces nothing', () => { + expect(keysOf([duty({ effective_from: '2026-09-01' })], now)).toEqual([]); + expect(keysOf([duty({ effective_to: '2026-07-31' })], now)).toEqual([]); + expect(planDispatch({ duties: [duty({ effective_to: '2026-07-31' })], now }).skipped[0]?.reason).toBe( + 'outside_effective_window', + ); + }); + + it('a recurring duty with no frequency is a FAULT, not a quiet skip', () => { + // `recurring_needs_frequency` should make this unreachable. If it is ever + // reached, the run must say so rather than silently dispatching nobody. + const plan = planDispatch({ duties: [duty({ frequency: null })], now }); + expect(plan.drafts).toEqual([]); + expect(plan.skipped[0]?.reason).toBe('no_frequency'); + }); + + it('a bad timezone stops that duty and nothing else', () => { + const plan = planDispatch({ duties: [duty({ id: 'bad', timezone: 'Mars/Olympus' }), duty({ id: 'good' })], now }); + expect(plan.skipped.map((s) => s.reason)).toEqual(['invalid_cadence']); + expect(plan.skipped[0]?.detail).toBeTruthy(); + expect(plan.drafts.map((d) => d.duty)).toEqual(['good']); + }); +}); + +describe('one UTC pass, every zone right', () => { + it('Asia/Shanghai and America/Los_Angeles get different local periods from one instant', () => { + // 2026-08-31T20:00Z is 2026-09-01 04:00 in Shanghai and 2026-08-31 13:00 in + // Los Angeles. Both duties are monthly; both answers are correct, and they + // are different months. + const now = new Date('2026-08-31T20:00:00Z'); + const plan = planDispatch({ + duties: [ + duty({ id: 'cn', timezone: 'Asia/Shanghai' }), + duty({ id: 'us', timezone: 'America/Los_Angeles' }), + ], + now, + }); + const byDuty = Object.fromEntries(plan.drafts.map((d) => [d.duty, d.period_key])); + expect(byDuty).toEqual({ cn: '2026-09', us: '2026-08' }); + }); + + it('the due date is resolved in the duty\'s zone too', () => { + const now = new Date('2026-08-31T20:00:00Z'); + const plan = planDispatch({ duties: [duty({ id: 'cn', timezone: 'Asia/Shanghai' })], now }); + // due_anchor period_start + 4 → the 5th of the duty's own September. + expect(plan.drafts[0]?.due_date).toBe('2026-09-05'); + }); + + it('every spelling in a draft is the period engine\'s, not this planner\'s', () => { + // The one assertion that makes the "sole authority" rule mean something: a + // draft is recomputed here straight from period.ts and must match byte for + // byte. A second derivation anywhere would show up as a mismatch. + const now = new Date('2026-08-15T09:00:00Z'); + for (const timezone of ['UTC', 'Asia/Shanghai', 'America/Los_Angeles', 'Pacific/Auckland']) { + for (const frequency of ['daily', 'weekly', 'fortnightly', 'monthly', 'quarterly', 'semiannual', 'annual'] as const) { + const d = duty({ frequency, timezone, lead_days: 3, due_anchor: 'period_end', due_offset_days: 0 }); + const [draft] = planDispatch({ duties: [d], now }).drafts; + const expectedKey = periodKeyFor(frequency, now, timezone); + const expectedDue = dueDateFor({ + frequency, + periodKey: expectedKey, + timezone, + dueAnchor: 'period_end', + dueOffsetDays: 0, + }); + const label = `${frequency} @ ${timezone}`; + expect(draft?.period_key, label).toBe(expectedKey); + expect(draft?.due_date, label).toBe(expectedDue); + expect(draft?.visible_from, label).toBe(visibleFromFor(expectedDue, 3)); + } + } + }); +}); + +describe('lead time decides how far ahead a task exists', () => { + it('with no lead time, only the current period', () => { + const now = new Date('2026-08-27T09:00:00Z'); + expect(keysOf([duty({ lead_days: 0 })], now)).toEqual(['2026-08']); + }); + + it('with lead time, next period\'s task exists once its window opens', () => { + // Monthly, due on the 3rd (period_start + 2), 7 days of lead → September's + // task becomes visible on 2026-08-27 and must therefore exist by then. + const d = duty({ due_offset_days: 2, lead_days: 7 }); + expect(keysOf([d], new Date('2026-08-26T09:00:00Z'))).toEqual(['2026-08']); + expect(keysOf([d], new Date('2026-08-27T09:00:00Z'))).toEqual(['2026-08', '2026-09']); + }); + + it('the current period exists even when its own lead window has not opened', () => { + // Due at period end with no lead: `visible_from` is 31 August, in the + // future on 1 August. `visible_from` governs when a task SHOWS UP, not + // whether it EXISTS, so the row is owed today. + const now = new Date('2026-08-01T09:00:00Z'); + const [draft] = planDispatch({ + duties: [duty({ due_anchor: 'period_end', due_offset_days: 0, lead_days: 0 })], + now, + }).drafts; + expect(draft?.period_key).toBe('2026-08'); + expect(draft?.due_date).toBe('2026-08-31'); + expect(draft?.visible_from).toBe('2026-08-31'); + }); + + it('the look-ahead is bounded by lead_days, so a long lead does not run away', () => { + // 400 days of lead on an annual duty reaches next year and stops there — + // not the year after, whose visible_from is still in the future. + const now = new Date('2026-08-15T09:00:00Z'); + const d = duty({ frequency: 'annual', due_anchor: 'period_start', due_offset_days: 0, lead_days: 400 }); + expect(keysOf([d], now)).toEqual(['2026', '2027']); + }); +}); + +describe('the copied fields', () => { + it('copies the subject and denormalises the business unit from the DUTY', () => { + const now = new Date('2026-08-15T09:00:00Z'); + const [draft] = planDispatch({ duties: [duty()], now }).drafts; + expect(draft).toMatchObject({ + subject: 'File the emissions return', + duty: 'duty_1', + owner: 'user_alice', + // From the duty, not the owner's current unit: a later transfer must not + // move historical rollups. + business_unit: 'bu_plant', + source: 'catalog', + status: 'open', + }); + }); + + it('never writes a status other than open, and never writes the server-owned stamps', () => { + const now = new Date('2026-08-15T09:00:00Z'); + const [draft] = planDispatch({ duties: [duty()], now }).drafts; + expect(draft?.status).toBe('open'); + // `completed_at` / `last_update_at` belong to task.hook.ts. A dispatcher + // that wrote either would be a second writer on the stagnation clock. + expect(Object.keys(draft ?? {})).not.toContain('completed_at'); + expect(Object.keys(draft ?? {})).not.toContain('last_update_at'); + }); +}); + +describe('backfill', () => { + const now = new Date('2026-08-15T09:00:00Z'); + + it('a 13-month window produces exactly the expected key set', () => { + const window = { from: '2025-08-01', to: '2026-08-31' }; + const keys = keysOf([duty()], now, window); + expect(keys).toEqual([ + '2025-08', '2025-09', '2025-10', '2025-11', '2025-12', + '2026-01', '2026-02', '2026-03', '2026-04', '2026-05', + '2026-06', '2026-07', '2026-08', + ]); + // And that set is `periodsBetween`'s, not a second walk written here. + expect(keys).toEqual( + periodsBetween('monthly', new Date('2025-08-01T12:00:00Z'), new Date('2026-08-31T12:00:00Z'), 'UTC'), + ); + }); + + it('clips to the duty\'s effective window, per period', () => { + const d = duty({ effective_from: '2026-03-10', effective_to: '2026-06-30' }); + // Due on the 5th, so March's due date (2026-03-05) predates the window. + expect(keysOf([d], now, { from: '2025-08-01', to: '2026-08-31' })).toEqual(['2026-04', '2026-05', '2026-06']); + }); + + it('backfills a duty whose window has already closed', () => { + // The duty-level "today is inside the window" gate is the SCHEDULED run's + // selection rule. Applying it to a backfill would make backfill useless for + // the case it exists for. + const d = duty({ effective_to: '2026-06-30' }); + expect(keysOf([d], now)).toEqual([]); + expect(keysOf([d], now, { from: '2026-04-01', to: '2026-08-31' })).toEqual(['2026-04', '2026-05', '2026-06']); + }); + + it('ignores the visibility gate — a past period is not "not visible yet"', () => { + const d = duty({ lead_days: 0, due_anchor: 'period_end' }); + expect(keysOf([d], now, { from: '2026-06-01', to: '2026-08-31' })).toEqual(['2026-06', '2026-07', '2026-08']); + }); +}); + +describe('parseBackfillWindow', () => { + const now = new Date('2026-08-15T09:00:00Z'); + + it('no input is the scheduled run', () => { + expect(parseBackfillWindow(undefined, now)).toBeNull(); + expect(parseBackfillWindow(null, now)).toBeNull(); + expect(parseBackfillWindow({}, now)).toBeNull(); + }); + + it('from alone runs up to the clock', () => { + expect(parseBackfillWindow({ from: '2026-01-01' }, now)).toEqual({ from: '2026-01-01', to: '2026-08-15' }); + }); + + it('both ends are honoured', () => { + expect(parseBackfillWindow({ from: '2026-01-01', to: '2026-03-31' }, now)).toEqual({ + from: '2026-01-01', + to: '2026-03-31', + }); + }); + + it('refuses a "to" with no floor rather than guessing one', () => { + expect(() => parseBackfillWindow({ to: '2026-03-31' }, now)).toThrow(/needs 'from'/); + }); + + it('refuses a malformed date and a reversed window', () => { + expect(() => parseBackfillWindow({ from: '01/01/2026' }, now)).toThrow(/YYYY-MM-DD/); + expect(() => parseBackfillWindow({ from: '2026-01-01', to: 7 }, now)).toThrow(/YYYY-MM-DD/); + expect(() => parseBackfillWindow({ from: '2026-05-01', to: '2026-01-01' }, now)).toThrow(/ends before it starts/); + expect(() => parseBackfillWindow('2026-01-01', now)).toThrow(/must be an object/); + }); +}); + +describe('nextDispatchedPeriod — advances, never regresses', () => { + it('takes the latest key created', () => { + expect(nextDispatchedPeriod(null, ['2026-08', '2026-09'])).toBe('2026-09'); + expect(nextDispatchedPeriod('2026-02', ['2026-08'])).toBe('2026-08'); + }); + + it('writes nothing when the run created nothing', () => { + expect(nextDispatchedPeriod('2026-08', [])).toBeNull(); + expect(nextDispatchedPeriod(null, [])).toBeNull(); + }); + + it('writes nothing rather than moving backwards', () => { + expect(nextDispatchedPeriod('2026-09', ['2026-08'])).toBeNull(); + expect(nextDispatchedPeriod('2026-08', ['2026-08'])).toBeNull(); + }); + + it('is chronological for every frequency spelling', () => { + expect(nextDispatchedPeriod('2026-W09', ['2026-W10'])).toBe('2026-W10'); + expect(nextDispatchedPeriod('2026-Q1', ['2026-Q3'])).toBe('2026-Q3'); + expect(nextDispatchedPeriod('2026-H1', ['2026-H2'])).toBe('2026-H2'); + expect(nextDispatchedPeriod('2025', ['2026'])).toBe('2026'); + expect(nextDispatchedPeriod('2026-08-09', ['2026-08-10'])).toBe('2026-08-10'); + // The zero padding is what makes a lexical compare chronological. + expect(nextDispatchedPeriod('2026-W09', ['2026-W08'])).toBeNull(); + }); +}); + +// ───────────────────────────────────────────────────────────────────────── +// 3. Idempotency, against a real engine that actually has the index +// ───────────────────────────────────────────────────────────────────────── + +let kernel: { getService(name: string): unknown; shutdown?(): Promise } | undefined; +/** + * The booted engine, typed as the rows it actually returns. + * + * Not `DispatchEngine`: that contract returns `unknown[]` so a caller cannot + * read a field it did not ask for, which is right for the dispatcher and + * useless for a test that has to inspect the rows. The real engine satisfies + * both, and the tests pass it straight to `runDispatch` structurally. + */ +let data: { + find(o: string, q?: AnyRow, x?: AnyRow): Promise; + insert(o: string, d: AnyRow, x?: AnyRow): Promise; + update(o: string, d: AnyRow, x?: AnyRow): Promise; +}; + +beforeAll(async () => { + const { plugins } = await createStandaloneStack({ + // sqlite, NOT memory: the memory driver stores no constraints, so the + // unique index this whole job rests on would not exist and every + // idempotency assertion below would pass without testing anything. + databaseDriver: 'sqlite', + databaseUrl: ':memory:', + skipSeedData: true, + // Left to its default this resolves `/dist/objectstack.json`, and a + // local `pnpm build` would make the suite report on the last BUILD rather + // than on `src/` — passing with the barrel entry deleted, and behaving + // differently in CI (where `pnpm test` runs before `pnpm build`). + artifactPath: 'dist/objectstack.this-suite-must-not-load-an-artifact.json', + }); + const k = new ObjectKernel(); + for (const plugin of plugins) await k.use(plugin); + await k.use(new AppPlugin(stack, undefined, { skipSeedData: true })); + await k.bootstrap(); + kernel = k as unknown as typeof kernel; + data = k.getService('data') as typeof data; +}, 180_000); + +afterAll(async () => { + await kernel?.shutdown?.(); +}); + +/** + * Every duty a test seeds, retired again when it ends. + * + * The suite shares one database, and `runDispatch` sweeps EVERY active + * recurring duty — so without this, a later test's run keeps creating tasks for + * an earlier test's duties and any assertion on a run-level count + * (`created === 0`) measures the whole file's history instead of this test. + * Retiring is enough: a retired duty leaves the sweep, and leaving the rows in + * place keeps each test's evidence readable if it fails. + */ +const seeded: string[] = []; + +afterEach(async () => { + unbindDispatchEngine(); + while (seeded.length > 0) { + const id = seeded.pop() as string; + await data.update('duly_duty', { status: 'retired' }, { where: { id }, multi: false }); + } +}); + +let seq = 0; +const seedDuty = async (over: AnyRow = {}): Promise => { + const created = await data.insert('duly_duty', { + name: `Duty ${++seq}`, + form: 'recurring', + owner: `user_${seq}`, + source: 'catalog', + status: 'active', + frequency: 'monthly', + due_anchor: 'period_start', + due_offset_days: 4, + lead_days: 0, + timezone: 'UTC', + ...over, + }); + const row = (Array.isArray(created) ? created[0] : created) as AnyRow; + seeded.push(String(row.id)); + return String(row.id); +}; + +const tasksFor = async (dutyId: string): Promise => + data.find('duly_task', { where: { duty: dutyId }, orderBy: [{ field: 'period_key', order: 'asc' }] }); + +const readDuty = async (dutyId: string): Promise => + (await data.find('duly_duty', { where: { id: dutyId }, limit: 1 }))[0] as AnyRow; + +const NOW = new Date('2026-08-15T09:00:00Z'); + +describe('the index is real on this engine', () => { + it('refuses a second (duty, owner, period_key) — which is what makes the job ordinary', async () => { + // Asserted rather than assumed. If this ever goes green-by-absence — a + // driver swap, an index rename — every idempotency test below becomes a + // test of nothing, and this is the one that says so. + const dutyId = await seedDuty(); + const row = { subject: 'x', duty: dutyId, owner: 'user_dup', period_key: '2099-01', source: 'catalog', status: 'open' }; + await data.insert('duly_task', row); + await expect(data.insert('duly_task', { ...row })).rejects.toThrow(); + expect(await data.find('duly_task', { where: { period_key: '2099-01' } })).toHaveLength(1); + }); +}); + +describe('running twice over the same clock', () => { + it('the second pass inserts nothing, and that is a successful run', async () => { + const dutyId = await seedDuty({ lead_days: 0 }); + + const first = await runDispatch(data, { now: NOW }); + expect(first.created).toBeGreaterThan(0); + const afterFirst = await tasksFor(dutyId); + expect(afterFirst.map((t) => t.period_key)).toEqual(['2026-08']); + + const second = await runDispatch(data, { now: NOW }); + expect(second.created).toBe(0); + expect(second.existing).toBeGreaterThan(0); + expect(second.degradedReason).toBeUndefined(); + expect(await tasksFor(dutyId)).toHaveLength(1); + }); + + it('two overlapping runs produce one task, not two', async () => { + // The race the unique index exists for. Both runs plan the same draft; one + // insert wins, the other fails, reads, finds the row and reports it as + // already existing. Neither run fails. + const dutyId = await seedDuty({ lead_days: 0 }); + const [a, b] = await Promise.all([runDispatch(data, { now: NOW }), runDispatch(data, { now: NOW })]); + expect(await tasksFor(dutyId)).toHaveLength(1); + const forThisDuty = (r: { created: number }) => r.created; + expect(forThisDuty(a) + forThisDuty(b)).toBeGreaterThan(0); + }); + + it('never touches an existing task row — the stagnation clock stays put', async () => { + // `last_update_at` is hook-stamped on every task write and deliberately + // does not advance on administrative ones. If the dispatcher ever updated + // a task, every stalled item would look freshly worked and the signal would + // go quiet with no error anywhere. + const dutyId = await seedDuty({ lead_days: 0 }); + await runDispatch(data, { now: NOW }); + const before = (await tasksFor(dutyId))[0]; + + const writes: string[] = []; + const spy: DispatchEngine = { + find: (o, q, x) => data.find(o, q as AnyRow, x as AnyRow), + insert: (o, d, x) => data.insert(o, d, x as AnyRow), + update: (o, d, x) => { + writes.push(o); + return data.update(o, d, x as AnyRow); + }, + }; + await runDispatch(spy, { now: NOW }); + await runDispatch(spy, { now: new Date('2026-09-15T09:00:00Z') }); + + expect(writes, 'the dispatcher must never update duly_task').not.toContain('duly_task'); + const after = (await tasksFor(dutyId))[0]; + expect(after.last_update_at).toBe(before.last_update_at); + expect(after.status).toBe('open'); + }); +}); + +describe('last_dispatched_period', () => { + it('advances to the latest key created, and stays put on a no-op run', async () => { + const dutyId = await seedDuty({ lead_days: 0 }); + await runDispatch(data, { now: NOW }); + expect((await readDuty(dutyId)).last_dispatched_period).toBe('2026-08'); + + const again = await runDispatch(data, { now: NOW }); + expect(again.advanced).toBe(0); + expect((await readDuty(dutyId)).last_dispatched_period).toBe('2026-08'); + + await runDispatch(data, { now: new Date('2026-09-15T09:00:00Z') }); + expect((await readDuty(dutyId)).last_dispatched_period).toBe('2026-09'); + }); + + it('never regresses when an older period is backfilled afterwards', async () => { + const dutyId = await seedDuty({ lead_days: 0 }); + await runDispatch(data, { now: NOW }); + expect((await readDuty(dutyId)).last_dispatched_period).toBe('2026-08'); + + await runDispatch(data, { now: NOW, window: { from: '2026-01-01', to: '2026-05-31' } }); + expect((await readDuty(dutyId)).last_dispatched_period).toBe('2026-08'); + expect((await tasksFor(dutyId)).map((t) => t.period_key)).toEqual([ + '2026-01', '2026-02', '2026-03', '2026-04', '2026-05', '2026-08', + ]); + }); +}); + +describe('backfill against the engine', () => { + it('a second backfill over the same window inserts nothing', async () => { + const dutyId = await seedDuty({ lead_days: 0, effective_from: '2025-01-01' }); + const window = { from: '2025-08-01', to: '2026-08-31' }; + + const first = await runDispatch(data, { now: NOW, window }); + expect(first.created).toBe(13); + expect((await tasksFor(dutyId))).toHaveLength(13); + + const second = await runDispatch(data, { now: NOW, window }); + expect(second.created).toBe(0); + expect(second.existing).toBe(13); + expect((await tasksFor(dutyId))).toHaveLength(13); + }); +}); + +describe('a real failure is re-raised, not swallowed as a duplicate', () => { + it('rethrows when the row is genuinely absent after a failed insert', async () => { + // The half a blanket try/catch gets wrong. A create that fails for a reason + // other than "already there" must fail the RUN — otherwise the dispatcher + // reports a clean night on which nothing was created. + await seedDuty({ lead_days: 0 }); + const boom = new Error('the store is on fire'); + const broken: DispatchEngine = { + find: (o, q, x) => data.find(o, q as AnyRow, x as AnyRow), + insert: (o) => (o === 'duly_task' ? Promise.reject(boom) : Promise.reject(boom)), + update: (o, d, x) => data.update(o, d, x as AnyRow), + }; + await expect(runDispatch(broken, { now: new Date('2027-03-15T09:00:00Z') })).rejects.toThrow('the store is on fire'); + }); +}); + +describe('paging', () => { + it('keeps reading while a page comes back full', async () => { + // The sweep must not stop at the first page. Driven with a fake rather than + // 201 real duties: the assertion is about the loop, not the store. + const pages: number[] = []; + const makeDuty = (i: number): DispatchDuty => ({ + id: `d${i}`, name: 'n', form: 'standing', status: 'active', owner: 'u', source: 'self', + }); + const paging: DispatchEngine = { + find: async (object, query): Promise => { + if (object !== 'duly_duty') return []; + const offset = Number((query as { offset?: number })?.offset ?? 0); + pages.push(offset); + if (offset === 0) return Array.from({ length: DISPATCH_PAGE_SIZE }, (_, i) => makeDuty(i)); + if (offset === DISPATCH_PAGE_SIZE) return [makeDuty(999)]; + return []; + }, + insert: async () => ({}), + update: async () => undefined, + }; + const result = await runDispatch(paging, { now: NOW }); + expect(pages).toEqual([0, DISPATCH_PAGE_SIZE]); + expect(result.duties).toBe(DISPATCH_PAGE_SIZE + 1); + }); +}); + +// ───────────────────────────────────────────────────────────────────────── +// 4. The handler, and the seam it needs +// ───────────────────────────────────────────────────────────────────────── + +describe('the job handler', () => { + it('refuses loudly when no host has bound an engine', async () => { + // The platform hands a job handler { jobId, data, bundle } and no engine. + // Until a host wires one, this must FAIL — a dispatcher that quietly does + // nothing is the worst failure this product can have. + unbindDispatchEngine(); + await expect(dulyDispatch({ jobId: DISPATCH_JOB_NAME })).rejects.toThrow(/no data engine/); + }); + + it('dispatches through the bound engine and reports completed', async () => { + const dutyId = await seedDuty({ lead_days: 0 }); + bindDispatchEngine(data); + const outcome = await dulyDispatch({ jobId: DISPATCH_JOB_NAME }); + expect(outcome).toEqual({ outcome: 'completed' }); + expect((await tasksFor(dutyId)).length).toBeGreaterThan(0); + }); + + it('takes a backfill window off the job input channel', async () => { + // `IJobService.trigger(name, data)` forwards `data` to the handler; this is + // the only non-schedule way into the job, and the reason the dispatcher is + // a job rather than a scheduled flow. + const dutyId = await seedDuty({ lead_days: 0, effective_from: '2026-01-01' }); + bindDispatchEngine(data); + await dulyDispatch({ jobId: DISPATCH_JOB_NAME, data: { from: '2026-03-01', to: '2026-05-31' } }); + expect((await tasksFor(dutyId)).map((t) => t.period_key)).toEqual(['2026-03', '2026-04', '2026-05']); + }); + + it('reports degraded — not failed — when a duty could not be dispatched', async () => { + // `degraded` is "ran to completion, work did not happen". It never retries, + // which is right: retrying a typo'd timezone at 01:05 will not fix it. + await seedDuty({ timezone: 'Mars/Olympus' }); + bindDispatchEngine(data); + const outcome = await dulyDispatch({ jobId: DISPATCH_JOB_NAME }); + expect(outcome.outcome).toBe('degraded'); + expect(outcome.reason).toMatch(/invalid_cadence/); + }); +}); + +describe('the product invariants this job must not undo', () => { + it('the identity index is still what dispatch relies on', () => { + const identity = Task.indexes?.find((i) => i.name === 'duly_task_dispatch_identity'); + expect(identity?.fields).toEqual(['duty', 'owner', 'period_key']); + expect(identity?.unique).toBe('organization'); + }); +});