This commit is contained in:
Daniel Hougaard
2024-01-22 23:54:21 +04:00
committed by Akhil Mohan
parent c21ea6fb75
commit edd78eaeba
3 changed files with 69 additions and 3 deletions

View File

@@ -112,6 +112,10 @@ export const queueServiceFactory = (redisUrl: string) => {
opts: JobsOptions & { jobId?: string } opts: JobsOptions & { jobId?: string }
) => { ) => {
const q = queueContainer[name]; const q = queueContainer[name];
console.log("name", name);
console.log("query", q);
await q.add(job, data, opts); await q.add(job, data, opts);
}; };

View File

@@ -350,7 +350,11 @@ export const registerRoutes = async (
integrationDAL, integrationDAL,
secretImportDAL, secretImportDAL,
projectEnvDAL, projectEnvDAL,
webhookDAL webhookDAL,
orgDAL,
projectMembershipDAL,
smtpService,
projectDAL
}); });
const secretService = secretServiceFactory({ const secretService = secretServiceFactory({

View File

@@ -10,11 +10,15 @@ import { QueueJobs, QueueName, TQueueServiceFactory } from "@app/queue";
import { TIntegrationDALFactory } from "../integration/integration-dal"; import { TIntegrationDALFactory } from "../integration/integration-dal";
import { TIntegrationAuthServiceFactory } from "../integration-auth/integration-auth-service"; import { TIntegrationAuthServiceFactory } from "../integration-auth/integration-auth-service";
import { syncIntegrationSecrets } from "../integration-auth/integration-sync-secret"; import { syncIntegrationSecrets } from "../integration-auth/integration-sync-secret";
import { TOrgDALFactory } from "../org/org-dal";
import { TProjectDALFactory } from "../project/project-dal";
import { TProjectBotServiceFactory } from "../project-bot/project-bot-service"; import { TProjectBotServiceFactory } from "../project-bot/project-bot-service";
import { TProjectEnvDALFactory } from "../project-env/project-env-dal"; import { TProjectEnvDALFactory } from "../project-env/project-env-dal";
import { TProjectMembershipDALFactory } from "../project-membership/project-membership-dal";
import { TSecretFolderDALFactory } from "../secret-folder/secret-folder-dal"; import { TSecretFolderDALFactory } from "../secret-folder/secret-folder-dal";
import { TSecretImportDALFactory } from "../secret-import/secret-import-dal"; import { TSecretImportDALFactory } from "../secret-import/secret-import-dal";
import { fnSecretsFromImports } from "../secret-import/secret-import-fns"; import { fnSecretsFromImports } from "../secret-import/secret-import-fns";
import { SmtpTemplates, TSmtpService } from "../smtp/smtp-service";
import { TWebhookDALFactory } from "../webhook/webhook-dal"; import { TWebhookDALFactory } from "../webhook/webhook-dal";
import { fnTriggerWebhook } from "../webhook/webhook-fns"; import { fnTriggerWebhook } from "../webhook/webhook-fns";
import { TSecretDALFactory } from "./secret-dal"; import { TSecretDALFactory } from "./secret-dal";
@@ -37,6 +41,10 @@ type TSecretQueueFactoryDep = {
secretImportDAL: Pick<TSecretImportDALFactory, "find">; secretImportDAL: Pick<TSecretImportDALFactory, "find">;
webhookDAL: Pick<TWebhookDALFactory, "findAllWebhooks" | "transaction" | "update" | "bulkUpdate">; webhookDAL: Pick<TWebhookDALFactory, "findAllWebhooks" | "transaction" | "update" | "bulkUpdate">;
projectEnvDAL: Pick<TProjectEnvDALFactory, "findOne">; projectEnvDAL: Pick<TProjectEnvDALFactory, "findOne">;
projectDAL: Pick<TProjectDALFactory, "findById">;
projectMembershipDAL: Pick<TProjectMembershipDALFactory, "findAllProjectMembers">;
smtpService: TSmtpService;
orgDAL: Pick<TOrgDALFactory, "findOrgByProjectId">;
}; };
export type TGetSecrets = { export type TGetSecrets = {
@@ -54,7 +62,11 @@ export const secretQueueFactory = ({
secretImportDAL, secretImportDAL,
folderDAL, folderDAL,
webhookDAL, webhookDAL,
projectEnvDAL projectEnvDAL,
orgDAL,
smtpService,
projectDAL,
projectMembershipDAL
}: TSecretQueueFactoryDep) => { }: TSecretQueueFactoryDep) => {
const syncIntegrations = async (dto: TGetSecrets) => { const syncIntegrations = async (dto: TGetSecrets) => {
await queueService.queue(QueueName.IntegrationSync, QueueJobs.IntegrationSync, dto, { await queueService.queue(QueueName.IntegrationSync, QueueJobs.IntegrationSync, dto, {
@@ -124,6 +136,8 @@ export const secretQueueFactory = ({
}); });
} }
console.log("Secret reminder thingies:)!");
console.log(queueService);
// If the secret already has a reminder, we should remove the existing one first. // If the secret already has a reminder, we should remove the existing one first.
if (oldSecret.secretReminderRepeatDays) { if (oldSecret.secretReminderRepeatDays) {
await removeSecretReminder({ await removeSecretReminder({
@@ -148,7 +162,8 @@ export const secretQueueFactory = ({
every: every:
appCfg.NODE_ENV === "development" appCfg.NODE_ENV === "development"
? secondsToMillis(newSecret.secretReminderRepeatDays) ? secondsToMillis(newSecret.secretReminderRepeatDays)
: daysToMillisecond(newSecret.secretReminderRepeatDays) : daysToMillisecond(newSecret.secretReminderRepeatDays),
immediately: true
} }
} }
); );
@@ -331,6 +346,49 @@ export const secretQueueFactory = ({
logger.info("Secret integration sync ended", job.id); logger.info("Secret integration sync ended", job.id);
}); });
queueService.start(QueueName.SecretReminder, async ({ data }) => {
logger.info(`secretReminderQueue.process: [secretDocument=${data.secretId}]`);
const { projectId } = data;
const organization = await orgDAL.findOrgByProjectId(projectId);
const project = await projectDAL.findById(projectId);
if (!organization) {
logger.info(
`secretReminderQueue.process: [secretDocument=${data.secretId}] no organization found`
);
return;
}
if (!project) {
logger.info(
`secretReminderQueue.process: [secretDocument=${data.secretId}] no project found`
);
return;
}
const projectMembers = await projectMembershipDAL.findAllProjectMembers(projectId);
if (!projectMembers || !projectMembers.length) {
logger.info(
`secretReminderQueue.process: [secretDocument=${data.secretId}] no project members found`
);
return;
}
await smtpService.sendMail({
template: SmtpTemplates.SecretReminder,
subjectLine: "Infisical secret reminder",
recipients: [...projectMembers.map((m) => m.user.email)],
substitutions: {
reminderNote: data.note, // May not be present.
projectName: project.name,
organizationName: organization.name
}
});
});
queueService.listen(QueueName.IntegrationSync, "failed", (job, err) => { queueService.listen(QueueName.IntegrationSync, "failed", (job, err) => {
logger.error("Failed to sync integration", job?.data, err); logger.error("Failed to sync integration", job?.data, err);
}); });