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