misc: finalized structure

This commit is contained in:
Sheen Capadngan
2024-12-04 22:49:28 +08:00
parent 11c96245a7
commit a750f48922
4 changed files with 100 additions and 93 deletions
+3
View File
@@ -10,12 +10,15 @@ export const mockQueue = (): TQueueServiceFactory => {
queue: async (name, jobData) => { queue: async (name, jobData) => {
job[name] = jobData; job[name] = jobData;
}, },
queuePg: async () => {},
initialize: async () => {},
shutdown: async () => undefined, shutdown: async () => undefined,
stopRepeatableJob: async () => true, stopRepeatableJob: async () => true,
start: (name, jobFn) => { start: (name, jobFn) => {
queues[name] = jobFn; queues[name] = jobFn;
workers[name] = jobFn; workers[name] = jobFn;
}, },
startPg: async () => {},
listen: (name, event) => { listen: (name, event) => {
events[name] = event; events[name] = event;
}, },
@@ -34,8 +34,9 @@ export const auditLogQueueServiceFactory = async ({
licenseService, licenseService,
auditLogStreamDAL auditLogStreamDAL
}: TAuditLogQueueServiceFactoryDep) => { }: TAuditLogQueueServiceFactoryDep) => {
const pushToLog = async (data: TCreateAuditLogDTO) => {
const appCfg = getConfig(); const appCfg = getConfig();
const pushToLog = async (data: TCreateAuditLogDTO) => {
if (appCfg.USE_PG_QUEUE) { if (appCfg.USE_PG_QUEUE) {
await queueService.queuePg<QueueName.AuditLog>(QueueJobs.AuditLog, data, { await queueService.queuePg<QueueName.AuditLog>(QueueJobs.AuditLog, data, {
retryLimit: 10, retryLimit: 10,
@@ -51,6 +52,7 @@ export const auditLogQueueServiceFactory = async ({
} }
}; };
if (appCfg.USE_PG_QUEUE) {
await queueService.startPg<QueueName.AuditLog>( await queueService.startPg<QueueName.AuditLog>(
QueueJobs.AuditLog, QueueJobs.AuditLog,
async ([job]) => { async ([job]) => {
@@ -141,6 +143,7 @@ export const auditLogQueueServiceFactory = async ({
pollingIntervalSeconds: 0.5 pollingIntervalSeconds: 0.5
} }
); );
}
queueService.start(QueueName.AuditLog, async (job) => { queueService.start(QueueName.AuditLog, async (job) => {
const { actor, event, ipAddress, projectId, userAgent, userAgentType } = job.data; const { actor, event, ipAddress, projectId, userAgent, userAgentType } = job.data;
-4
View File
@@ -57,11 +57,7 @@ const run = async () => {
const smtp = smtpServiceFactory(formatSmtpConfig()); const smtp = smtpServiceFactory(formatSmtpConfig());
const queue = queueServiceFactory(appCfg.REDIS_URL, appCfg.DB_CONNECTION_URI); const queue = queueServiceFactory(appCfg.REDIS_URL, appCfg.DB_CONNECTION_URI);
if (appCfg.USE_PG_QUEUE) {
logger.info("Initializing PG queue...");
await queue.initialize(); await queue.initialize();
}
const keyStore = keyStoreFactory(appCfg.REDIS_URL); const keyStore = keyStoreFactory(appCfg.REDIS_URL);
+7 -2
View File
@@ -8,6 +8,7 @@ import {
TScanFullRepoEventPayload, TScanFullRepoEventPayload,
TScanPushEventPayload TScanPushEventPayload
} from "@app/ee/services/secret-scanning/secret-scanning-queue/secret-scanning-queue-types"; } from "@app/ee/services/secret-scanning/secret-scanning-queue/secret-scanning-queue-types";
import { getConfig } from "@app/lib/config/env";
import { logger } from "@app/lib/logger"; import { logger } from "@app/lib/logger";
import { import {
TFailedIntegrationSyncEmailsPayload, TFailedIntegrationSyncEmailsPayload,
@@ -195,8 +196,8 @@ export const queueServiceFactory = (redisUrl: string, dbConnectionUrl: string) =
const pgBoss = new PgBoss({ const pgBoss = new PgBoss({
connectionString: dbConnectionUrl, connectionString: dbConnectionUrl,
archiveCompletedAfterSeconds: 30, archiveCompletedAfterSeconds: 60,
maintenanceIntervalSeconds: 30, archiveFailedAfterSeconds: 1000, // we want to keep failed jobs for a longer time so that it can be retried
deleteAfterSeconds: 30 deleteAfterSeconds: 30
}); });
@@ -208,11 +209,15 @@ export const queueServiceFactory = (redisUrl: string, dbConnectionUrl: string) =
>; >;
const initialize = async () => { const initialize = async () => {
const appCfg = getConfig();
if (appCfg.USE_PG_QUEUE) {
logger.info("Initializing PG queue...");
await pgBoss.start(); await pgBoss.start();
pgBoss.on("error", (error) => { pgBoss.on("error", (error) => {
logger.error(error, "pg-queue error"); logger.error(error, "pg-queue error");
}); });
}
}; };
const start = <T extends QueueName>( const start = <T extends QueueName>(