feat(12-02): TenderDigestScheduler — ein globaler Cron, findMany über fällige Nutzer
A single platform-wide @nestjs/schedule cron (daily 07:00), registered via SchedulerRegistry exactly like TenderSchedulerService — NOT the DkvSchedulerService single-tenant pattern (Pitfall 1). Selects candidate users as distinct userId with an open TenderMatch (notifiedAt IS NULL) via findMany across all tenants, resolves each user's TenderNotificationPref.digestInterval (missing row -> daily default, D-01: daily always due, weekly only on Monday Europe/Berlin, off never), groups their un-notified matches by saved-search profile name into one TenderMailService.sendDigest call per user (D-02), and stamps notifiedAt+channel='digest' ONLY after a successful send — the shared notifiedAt-IS-NULL eligibility gate that guarantees no double-send with instant alerts (D-06). Each candidate user is processed in its own try/catch: a missing SMTP config, a send failure, or an unexpected thrown error for one user/tenant leaves that user's matches notifiedAt=NULL (retried next run) and never aborts the run for the rest (Pitfall 6). Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,186 @@
|
||||
import { Injectable, Logger, OnModuleInit } from '@nestjs/common';
|
||||
import { SchedulerRegistry } from '@nestjs/schedule';
|
||||
import { PrismaService } from '../prisma/prisma.service';
|
||||
import { TenderMailItem, TenderMailService } from './tender-mail.service';
|
||||
|
||||
/**
|
||||
* CronJob constructor — resolved at runtime via require() because `cron` is
|
||||
* a transitive dependency of @nestjs/schedule (not a direct api dep under
|
||||
* pnpm strict isolation, so `import { CronJob } from 'cron'` fails
|
||||
* type-check). Reuses the exact `TenderSchedulerService`/`DkvSchedulerService`
|
||||
* resolution workaround verbatim.
|
||||
*/
|
||||
// eslint-disable-next-line @typescript-eslint/no-require-imports
|
||||
const CronJobClass: new (cronTime: string, onTick: () => void) => { start(): void } =
|
||||
// eslint-disable-next-line @typescript-eslint/no-unsafe-member-access
|
||||
require('cron').CronJob as new (cronTime: string, onTick: () => void) => { start(): void };
|
||||
|
||||
/**
|
||||
* TenderDigestScheduler — a SINGLE global cron job driving the tender-radar
|
||||
* digest send (NOTIFY-01). Deliberately mirrors `TenderSchedulerService`'s
|
||||
* poll-once-fan-out-many pattern, NOT `DkvSchedulerService`'s documented
|
||||
* "v1 single-tenant, find-first-row" pattern (RESEARCH.md Pitfall 1): there is
|
||||
* no per-tenant cron job and no per-tenant "active tenant" instance field.
|
||||
* One job, registered once, iterates every due user of every tenant via
|
||||
* `findMany` on every tick.
|
||||
*
|
||||
* Runs daily at 07:00 (server-local cron tick; due-date evaluation itself
|
||||
* is Europe/Berlin-aware for the weekly weekday check). Selects candidate
|
||||
* users as the distinct `userId`s that have at least one un-notified
|
||||
* `TenderMatch`, resolves each user's `TenderNotificationPref.digestInterval`
|
||||
* (missing row -> 'daily' default, D-01), and — for due users only — groups
|
||||
* their un-notified matches by saved-search profile name (D-02) into ONE
|
||||
* `TenderMailService.sendDigest` call. `TenderMatch.notifiedAt`/`notifiedChannel`
|
||||
* are stamped ONLY after a successful (non-skipped) send — the single
|
||||
* `notifiedAt IS NULL` eligibility gate that structurally prevents
|
||||
* double-sends across digest and instant channels (D-06).
|
||||
*
|
||||
* Robustness (Pitfall 6): each candidate user is processed inside its own
|
||||
* try/catch. A missing SMTP config, a send failure, or an unexpected thrown
|
||||
* error for one user/tenant leaves that user's matches `notifiedAt=NULL`
|
||||
* (retried on the next run) and never aborts the run for the remaining
|
||||
* users — a broken tenant must never take down every other tenant's digest.
|
||||
*/
|
||||
@Injectable()
|
||||
export class TenderDigestScheduler implements OnModuleInit {
|
||||
private readonly logger = new Logger(TenderDigestScheduler.name);
|
||||
|
||||
/** Name of the single, platform-global managed cron job. */
|
||||
private readonly JOB_NAME = 'tender-digest';
|
||||
|
||||
/** Daily at 07:00 — a single platform-wide tick, no tenant dimension. */
|
||||
private readonly CRON_EXPR = '0 7 * * *';
|
||||
|
||||
constructor(
|
||||
private readonly schedulerRegistry: SchedulerRegistry,
|
||||
private readonly prisma: PrismaService,
|
||||
private readonly mail: TenderMailService,
|
||||
) {}
|
||||
|
||||
/**
|
||||
* Registers the single global digest cron job on application startup.
|
||||
* Errors are caught and logged (never re-thrown) so a scheduling issue
|
||||
* never prevents the rest of the application from starting.
|
||||
*/
|
||||
onModuleInit(): void {
|
||||
try {
|
||||
// Remove existing job if already registered (e.g. hot-reload/tests).
|
||||
try {
|
||||
this.schedulerRegistry.getCronJob(this.JOB_NAME).stop();
|
||||
this.schedulerRegistry.deleteCronJob(this.JOB_NAME);
|
||||
} catch {
|
||||
/* Job not yet registered — expected on first boot */
|
||||
}
|
||||
|
||||
const job = new CronJobClass(this.CRON_EXPR, () => {
|
||||
this.runDigest().catch((err) =>
|
||||
this.logger.error(`Tender digest run failed: ${(err as Error).message}`),
|
||||
);
|
||||
});
|
||||
|
||||
// Cast required: our minimal CronJob type doesn't match cron's full
|
||||
// type signature. At runtime the object IS a full CronJob —
|
||||
// SchedulerRegistry only calls stop() on it.
|
||||
// eslint-disable-next-line @typescript-eslint/no-explicit-any
|
||||
this.schedulerRegistry.addCronJob(this.JOB_NAME, job as any);
|
||||
job.start();
|
||||
|
||||
this.logger.log(
|
||||
`Tender digest scheduler registered: ${this.CRON_EXPR} (single global job — all tenants/users, no per-tenant dimension)`,
|
||||
);
|
||||
} catch (err) {
|
||||
this.logger.error(`Tender digest scheduler init failed: ${(err as Error).message}`);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Core digest run — the single global tick's fan-out over every due user
|
||||
* of every tenant. NEVER a per-tenant job, NEVER a first-row-only lookup (Pitfall 1).
|
||||
*
|
||||
* @param now - injectable clock for the weekly-weekday check (defaults to
|
||||
* the real current time); tests pass a fixed date instead of faking the
|
||||
* system clock.
|
||||
*/
|
||||
async runDigest(now: Date = new Date()): Promise<void> {
|
||||
// Candidate users: distinct userId with at least one un-notified match,
|
||||
// across ALL tenants — a single findMany, never a per-tenant iteration.
|
||||
const candidates = await this.prisma.tenderMatch.findMany({
|
||||
where: { notifiedAt: null },
|
||||
select: { userId: true },
|
||||
distinct: ['userId'],
|
||||
});
|
||||
|
||||
if (!candidates.length) return;
|
||||
|
||||
const weeklyDue = isMondayInBerlin(now);
|
||||
|
||||
for (const { userId } of candidates) {
|
||||
try {
|
||||
const pref = await this.prisma.tenderNotificationPref.findUnique({
|
||||
where: { userId },
|
||||
});
|
||||
// Missing pref row -> daily default (D-01).
|
||||
const interval = pref?.digestInterval ?? 'daily';
|
||||
|
||||
if (interval === 'off') continue;
|
||||
if (interval === 'weekly' && !weeklyDue) continue;
|
||||
// 'daily' (or any unrecognized value) -> due every run.
|
||||
|
||||
const matches = await this.prisma.tenderMatch.findMany({
|
||||
where: { userId, notifiedAt: null },
|
||||
include: { tender: true, savedSearch: true },
|
||||
orderBy: { savedSearch: { name: 'asc' } },
|
||||
});
|
||||
if (!matches.length) continue;
|
||||
|
||||
const user = await this.prisma.user.findUnique({ where: { id: userId } });
|
||||
if (!user) continue;
|
||||
|
||||
const sections = groupMatchesByProfile(matches);
|
||||
const sent = await this.mail.sendDigest({ email: user.email }, user.tenantId, sections);
|
||||
|
||||
if (sent) {
|
||||
await this.prisma.tenderMatch.updateMany({
|
||||
where: { id: { in: matches.map((m: { id: string }) => m.id) } },
|
||||
data: { notifiedAt: new Date(), notifiedChannel: 'digest' },
|
||||
});
|
||||
}
|
||||
// sent === false (no SMTP config or send failure) -> notifiedAt
|
||||
// stays NULL, this user's matches are retried on the next run.
|
||||
} catch (err) {
|
||||
// One broken user/tenant (DB error, malformed pref, etc.) must
|
||||
// never abort the run for the remaining users (Pitfall 6).
|
||||
this.logger.error(
|
||||
`Tender digest failed for user ${userId}: ${(err as Error).message}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Groups a flat list of matches (each including its `savedSearch` and
|
||||
* `tender` relations) into the `{ profileName: tender[] }` shape
|
||||
* `TenderMailService.sendDigest` expects (D-02 — one mail, sectioned by
|
||||
* saved-search profile).
|
||||
*/
|
||||
function groupMatchesByProfile(
|
||||
matches: Array<{ savedSearch: { name: string }; tender: TenderMailItem }>,
|
||||
): Record<string, TenderMailItem[]> {
|
||||
const sections: Record<string, TenderMailItem[]> = {};
|
||||
for (const match of matches) {
|
||||
const profileName = match.savedSearch.name;
|
||||
if (!sections[profileName]) sections[profileName] = [];
|
||||
sections[profileName].push(match.tender);
|
||||
}
|
||||
return sections;
|
||||
}
|
||||
|
||||
/** True when `now` falls on a Monday in the Europe/Berlin timezone. */
|
||||
function isMondayInBerlin(now: Date): boolean {
|
||||
const weekday = now.toLocaleDateString('en-US', {
|
||||
timeZone: 'Europe/Berlin',
|
||||
weekday: 'short',
|
||||
});
|
||||
return weekday === 'Mon';
|
||||
}
|
||||
Reference in New Issue
Block a user