Press n or j to go to the next uncovered block, b, p or k for the previous block.
| 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 | 60x 60x 60x 60x 60x 60x 60x 60x 139x 139x 109x 109x 109x 139x 81x 243x 81x 81x 139x 134x 221x 31x 205x 200x 489x 200x 489x 489x 489x 139x 20x 24x 11x 60x 5x 97x | import { queueAsPromised, promise as queuePromise } from "fastq";
import { QueueError } from "shared/types/api/errors";
import { ProgramType } from "shared/types/models/programItem";
import {
Result,
makeErrorResult,
makeSuccessResult,
} from "shared/utils/result";
import { EmailSender } from "server/features/notifications/email";
import {
EmailNotificationOutcome,
emailNotificationWorker,
} from "server/features/notifications/emailNotificationWorker";
import { logger } from "server/utils/logger";
export enum NotificationTaskType {
SEND_EMAIL_ACCEPTED,
SEND_EMAIL_REJECTED,
SEND_EMAIL_PROGRAM_ITEM_CANCELLED,
SEND_EMAIL_PROGRAM_ITEM_DELETED,
SEND_EMAIL_PROGRAM_ITEM_NO_KONSTI_SIGNUP_ANYMORE,
SEND_EMAIL_PROGRAM_ITEM_NO_LOTTERY_ANYMORE,
SEND_EMAIL_PROGRAM_ITEM_TIME_CHANGED,
}
export interface NotificationTask {
type: NotificationTaskType;
username: string;
programItemId: string;
programItemStartTime: string;
// The end of the last program item a batched lottery covered, paired with the start time above
// so the rejection can name the whole span, and the program type to name what was lotteried
lastProgramItemEndTime?: string;
programType?: ProgramType;
programItemTitle?: string;
}
export interface NotificationQueueService {
addNotificationsBulk(
notifications: NotificationTask[],
): Result<boolean, QueueError>;
drain(): Promise<void>;
kill(): Promise<void>;
getItems(): NotificationTask[];
getQueue(): queueAsPromised<NotificationTask, EmailNotificationOutcome>;
getSender(): EmailSender;
}
export function createNotificationQueueService(
sender: EmailSender,
workerCount = 1,
stopOnStart = false,
): NotificationQueueService {
// The worker logs each failure on its own, so this is the one line saying how a batch went.
// Tallied inside the worker rather than off the push promise, because the queue calls drain
// before the last task's promise settles.
const outcomes = new Map<EmailNotificationOutcome, number>();
const queue: queueAsPromised<NotificationTask, EmailNotificationOutcome> =
queuePromise(async (notification: NotificationTask) => {
const outcome = await emailNotificationWorker(sender, notification);
outcomes.set(outcome, (outcomes.get(outcome) ?? 0) + 1);
return outcome;
}, workerCount);
// fastq calls this hook each time the last running task finishes with nothing waiting, so it
// runs once per batch, after all its emails have been processed
queue.drain = () => {
const count = (outcome: EmailNotificationOutcome): number =>
outcomes.get(outcome) ?? 0;
logger.info(
`Email notification queue drained: ${count(EmailNotificationOutcome.SENT)} sent, ${count(EmailNotificationOutcome.SKIPPED)} skipped (no email address), ${count(EmailNotificationOutcome.FAILED)} failed`,
);
outcomes.clear();
};
Eif (stopOnStart) {
queue.pause();
}
function addNotificationsBulk(
notifications: NotificationTask[],
): Result<boolean, QueueError> {
if (notifications.length === 0) {
return makeSuccessResult(true);
}
try {
for (const notification of notifications) {
addNotification(notification);
}
return makeSuccessResult(true);
} catch {
return makeErrorResult(QueueError.FAILED_TO_PUSH);
}
}
function addNotification(
notification: NotificationTask,
): Result<boolean, QueueError> {
try {
// The promise push returns settles only once the task has run, and nothing waits for that
void queue.push(notification);
return makeSuccessResult(true);
} catch {
return makeErrorResult(QueueError.FAILED_TO_PUSH);
}
}
return {
addNotificationsBulk,
drain: async () => {
await queue.drain();
},
kill: async () => {
await queue.kill();
},
getItems(): NotificationTask[] {
return queue.getQueue();
},
getQueue(): queueAsPromised<NotificationTask, EmailNotificationOutcome> {
return queue;
},
getSender(): EmailSender {
return sender;
},
};
}
let globalNotificationQueueService: NotificationQueueService | null = null;
export function setGlobalNotificationQueueService(
service: NotificationQueueService | null,
): void {
globalNotificationQueueService = service;
}
export function getGlobalNotificationQueueService(): NotificationQueueService | null {
return globalNotificationQueueService;
}
|