mirror of
https://github.com/awatertrevi/infisical.git
synced 2026-10-06 18:27:19 +00:00
queue job
This commit is contained in:
Vendored
+1
-1
@@ -71,6 +71,7 @@ import { TIdentityTokenAuthServiceFactory } from "@app/services/identity-token-a
|
|||||||
import { TIdentityUaServiceFactory } from "@app/services/identity-ua/identity-ua-service";
|
import { TIdentityUaServiceFactory } from "@app/services/identity-ua/identity-ua-service";
|
||||||
import { TIntegrationServiceFactory } from "@app/services/integration/integration-service";
|
import { TIntegrationServiceFactory } from "@app/services/integration/integration-service";
|
||||||
import { TIntegrationAuthServiceFactory } from "@app/services/integration-auth/integration-auth-service";
|
import { TIntegrationAuthServiceFactory } from "@app/services/integration-auth/integration-auth-service";
|
||||||
|
import { TMicrosoftTeamsServiceFactory } from "@app/services/microsoft-teams/microsoft-teams-service";
|
||||||
import { TOrgRoleServiceFactory } from "@app/services/org/org-role-service";
|
import { TOrgRoleServiceFactory } from "@app/services/org/org-role-service";
|
||||||
import { TOrgServiceFactory } from "@app/services/org/org-service";
|
import { TOrgServiceFactory } from "@app/services/org/org-service";
|
||||||
import { TOrgAdminServiceFactory } from "@app/services/org-admin/org-admin-service";
|
import { TOrgAdminServiceFactory } from "@app/services/org-admin/org-admin-service";
|
||||||
@@ -100,7 +101,6 @@ import { TUserServiceFactory } from "@app/services/user/user-service";
|
|||||||
import { TUserEngagementServiceFactory } from "@app/services/user-engagement/user-engagement-service";
|
import { TUserEngagementServiceFactory } from "@app/services/user-engagement/user-engagement-service";
|
||||||
import { TWebhookServiceFactory } from "@app/services/webhook/webhook-service";
|
import { TWebhookServiceFactory } from "@app/services/webhook/webhook-service";
|
||||||
import { TWorkflowIntegrationServiceFactory } from "@app/services/workflow-integration/workflow-integration-service";
|
import { TWorkflowIntegrationServiceFactory } from "@app/services/workflow-integration/workflow-integration-service";
|
||||||
import { TMicrosoftTeamsServiceFactory } from "@app/services/microsoft-teams/microsoft-teams-service";
|
|
||||||
|
|
||||||
declare module "@fastify/request-context" {
|
declare module "@fastify/request-context" {
|
||||||
interface RequestContextData {
|
interface RequestContextData {
|
||||||
|
|||||||
@@ -25,6 +25,7 @@ import {
|
|||||||
TQueueSecretSyncSyncSecretsByIdDTO,
|
TQueueSecretSyncSyncSecretsByIdDTO,
|
||||||
TQueueSendSecretSyncActionFailedNotificationsDTO
|
TQueueSendSecretSyncActionFailedNotificationsDTO
|
||||||
} from "@app/services/secret-sync/secret-sync-types";
|
} from "@app/services/secret-sync/secret-sync-types";
|
||||||
|
import { CacheType } from "@app/services/super-admin/super-admin-types";
|
||||||
import { TWebhookPayloads } from "@app/services/webhook/webhook-types";
|
import { TWebhookPayloads } from "@app/services/webhook/webhook-types";
|
||||||
|
|
||||||
export enum QueueName {
|
export enum QueueName {
|
||||||
@@ -49,7 +50,8 @@ export enum QueueName {
|
|||||||
AccessTokenStatusUpdate = "access-token-status-update",
|
AccessTokenStatusUpdate = "access-token-status-update",
|
||||||
ImportSecretsFromExternalSource = "import-secrets-from-external-source",
|
ImportSecretsFromExternalSource = "import-secrets-from-external-source",
|
||||||
AppConnectionSecretSync = "app-connection-secret-sync",
|
AppConnectionSecretSync = "app-connection-secret-sync",
|
||||||
SecretRotationV2 = "secret-rotation-v2"
|
SecretRotationV2 = "secret-rotation-v2",
|
||||||
|
InvalidateCache = "invalidate-cache"
|
||||||
}
|
}
|
||||||
|
|
||||||
export enum QueueJobs {
|
export enum QueueJobs {
|
||||||
@@ -81,7 +83,8 @@ export enum QueueJobs {
|
|||||||
SecretSyncSendActionFailedNotifications = "secret-sync-send-action-failed-notifications",
|
SecretSyncSendActionFailedNotifications = "secret-sync-send-action-failed-notifications",
|
||||||
SecretRotationV2QueueRotations = "secret-rotation-v2-queue-rotations",
|
SecretRotationV2QueueRotations = "secret-rotation-v2-queue-rotations",
|
||||||
SecretRotationV2RotateSecrets = "secret-rotation-v2-rotate-secrets",
|
SecretRotationV2RotateSecrets = "secret-rotation-v2-rotate-secrets",
|
||||||
SecretRotationV2SendNotification = "secret-rotation-v2-send-notification"
|
SecretRotationV2SendNotification = "secret-rotation-v2-send-notification",
|
||||||
|
InvalidateCache = "invalidate-cache"
|
||||||
}
|
}
|
||||||
|
|
||||||
export type TQueueJobTypes = {
|
export type TQueueJobTypes = {
|
||||||
@@ -234,6 +237,14 @@ export type TQueueJobTypes = {
|
|||||||
name: QueueJobs.SecretRotationV2SendNotification;
|
name: QueueJobs.SecretRotationV2SendNotification;
|
||||||
payload: TSecretRotationSendNotificationJobPayload;
|
payload: TSecretRotationSendNotificationJobPayload;
|
||||||
};
|
};
|
||||||
|
[QueueName.InvalidateCache]: {
|
||||||
|
name: QueueJobs.InvalidateCache;
|
||||||
|
payload: {
|
||||||
|
data: {
|
||||||
|
type: CacheType;
|
||||||
|
};
|
||||||
|
};
|
||||||
|
};
|
||||||
};
|
};
|
||||||
|
|
||||||
export type TQueueServiceFactory = ReturnType<typeof queueServiceFactory>;
|
export type TQueueServiceFactory = ReturnType<typeof queueServiceFactory>;
|
||||||
|
|||||||
@@ -238,6 +238,7 @@ import { projectSlackConfigDALFactory } from "@app/services/slack/project-slack-
|
|||||||
import { slackIntegrationDALFactory } from "@app/services/slack/slack-integration-dal";
|
import { slackIntegrationDALFactory } from "@app/services/slack/slack-integration-dal";
|
||||||
import { slackServiceFactory } from "@app/services/slack/slack-service";
|
import { slackServiceFactory } from "@app/services/slack/slack-service";
|
||||||
import { TSmtpService } from "@app/services/smtp/smtp-service";
|
import { TSmtpService } from "@app/services/smtp/smtp-service";
|
||||||
|
import { invalidateCacheQueueFactory } from "@app/services/super-admin/invalidate-cache-queue";
|
||||||
import { superAdminDALFactory } from "@app/services/super-admin/super-admin-dal";
|
import { superAdminDALFactory } from "@app/services/super-admin/super-admin-dal";
|
||||||
import { getServerCfg, superAdminServiceFactory } from "@app/services/super-admin/super-admin-service";
|
import { getServerCfg, superAdminServiceFactory } from "@app/services/super-admin/super-admin-service";
|
||||||
import { telemetryDALFactory } from "@app/services/telemetry/telemetry-dal";
|
import { telemetryDALFactory } from "@app/services/telemetry/telemetry-dal";
|
||||||
@@ -605,6 +606,11 @@ export const registerRoutes = async (
|
|||||||
queueService
|
queueService
|
||||||
});
|
});
|
||||||
|
|
||||||
|
const invalidateCacheQueue = invalidateCacheQueueFactory({
|
||||||
|
keyStore,
|
||||||
|
queueService
|
||||||
|
});
|
||||||
|
|
||||||
const userService = userServiceFactory({
|
const userService = userServiceFactory({
|
||||||
userDAL,
|
userDAL,
|
||||||
userAliasDAL,
|
userAliasDAL,
|
||||||
@@ -715,7 +721,8 @@ export const registerRoutes = async (
|
|||||||
keyStore,
|
keyStore,
|
||||||
licenseService,
|
licenseService,
|
||||||
kmsService,
|
kmsService,
|
||||||
microsoftTeamsService
|
microsoftTeamsService,
|
||||||
|
invalidateCacheQueue
|
||||||
});
|
});
|
||||||
|
|
||||||
const orgAdminService = orgAdminServiceFactory({
|
const orgAdminService = orgAdminServiceFactory({
|
||||||
|
|||||||
@@ -572,19 +572,18 @@ export const registerAdminRouter = async (server: FastifyZodProvider) => {
|
|||||||
});
|
});
|
||||||
},
|
},
|
||||||
handler: async (req) => {
|
handler: async (req) => {
|
||||||
const keysCleared = await server.services.superAdmin.invalidateCache(req.body.type);
|
await server.services.superAdmin.invalidateCache(req.body.type);
|
||||||
|
|
||||||
await server.services.telemetry.sendPostHogEvents({
|
await server.services.telemetry.sendPostHogEvents({
|
||||||
event: PostHogEventTypes.InvalidateCache,
|
event: PostHogEventTypes.InvalidateCache,
|
||||||
distinctId: getTelemetryDistinctId(req),
|
distinctId: getTelemetryDistinctId(req),
|
||||||
properties: {
|
properties: {
|
||||||
keysCleared,
|
|
||||||
...req.auditLogInfo
|
...req.auditLogInfo
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
return {
|
return {
|
||||||
message: `Successfully invalidated ${keysCleared} cached items`
|
message: "Cache invalidation job started"
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -26,6 +26,7 @@ import { registerIdentityUaRouter } from "./identity-universal-auth-router";
|
|||||||
import { registerIntegrationAuthRouter } from "./integration-auth-router";
|
import { registerIntegrationAuthRouter } from "./integration-auth-router";
|
||||||
import { registerIntegrationRouter } from "./integration-router";
|
import { registerIntegrationRouter } from "./integration-router";
|
||||||
import { registerInviteOrgRouter } from "./invite-org-router";
|
import { registerInviteOrgRouter } from "./invite-org-router";
|
||||||
|
import { registerMicrosoftTeamsRouter } from "./microsoft-teams-router";
|
||||||
import { registerOrgAdminRouter } from "./org-admin-router";
|
import { registerOrgAdminRouter } from "./org-admin-router";
|
||||||
import { registerOrgRouter } from "./organization-router";
|
import { registerOrgRouter } from "./organization-router";
|
||||||
import { registerPasswordRouter } from "./password-router";
|
import { registerPasswordRouter } from "./password-router";
|
||||||
@@ -47,7 +48,6 @@ import { registerUserEngagementRouter } from "./user-engagement-router";
|
|||||||
import { registerUserRouter } from "./user-router";
|
import { registerUserRouter } from "./user-router";
|
||||||
import { registerWebhookRouter } from "./webhook-router";
|
import { registerWebhookRouter } from "./webhook-router";
|
||||||
import { registerWorkflowIntegrationRouter } from "./workflow-integration-router";
|
import { registerWorkflowIntegrationRouter } from "./workflow-integration-router";
|
||||||
import { registerMicrosoftTeamsRouter } from "./microsoft-teams-router";
|
|
||||||
|
|
||||||
export const registerV1Routes = async (server: FastifyZodProvider) => {
|
export const registerV1Routes = async (server: FastifyZodProvider) => {
|
||||||
await server.register(registerSsoRouter, { prefix: "/sso" });
|
await server.register(registerSsoRouter, { prefix: "/sso" });
|
||||||
|
|||||||
@@ -0,0 +1,43 @@
|
|||||||
|
import { TKeyStoreFactory } from "@app/keystore/keystore";
|
||||||
|
import { logger } from "@app/lib/logger";
|
||||||
|
import { QueueJobs, QueueName, TQueueServiceFactory } from "@app/queue";
|
||||||
|
|
||||||
|
import { CacheType } from "./super-admin-types";
|
||||||
|
|
||||||
|
export type TInvalidateCacheQueueFactoryDep = {
|
||||||
|
queueService: TQueueServiceFactory;
|
||||||
|
|
||||||
|
keyStore: Pick<TKeyStoreFactory, "deleteItems">;
|
||||||
|
};
|
||||||
|
|
||||||
|
export type TInvalidateCacheQueueFactory = ReturnType<typeof invalidateCacheQueueFactory>;
|
||||||
|
|
||||||
|
export const invalidateCacheQueueFactory = ({ queueService, keyStore }: TInvalidateCacheQueueFactoryDep) => {
|
||||||
|
const startInvalidate = async (dto: {
|
||||||
|
data: {
|
||||||
|
type: CacheType;
|
||||||
|
};
|
||||||
|
}) => {
|
||||||
|
await queueService.queue(QueueName.InvalidateCache, QueueJobs.InvalidateCache, dto, {
|
||||||
|
removeOnComplete: true,
|
||||||
|
removeOnFail: true
|
||||||
|
});
|
||||||
|
};
|
||||||
|
|
||||||
|
queueService.start(QueueName.InvalidateCache, async (job) => {
|
||||||
|
try {
|
||||||
|
const {
|
||||||
|
data: { type }
|
||||||
|
} = job.data;
|
||||||
|
|
||||||
|
if (type === CacheType.ALL || type === CacheType.SECRETS)
|
||||||
|
await keyStore.deleteItems({ pattern: "secret-manager:*" });
|
||||||
|
} catch (err) {
|
||||||
|
logger.error(err, "Failed to invalidate cache");
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
return {
|
||||||
|
startInvalidate
|
||||||
|
};
|
||||||
|
};
|
||||||
@@ -25,6 +25,7 @@ import { TOrgServiceFactory } from "../org/org-service";
|
|||||||
import { TUserDALFactory } from "../user/user-dal";
|
import { TUserDALFactory } from "../user/user-dal";
|
||||||
import { TUserAliasDALFactory } from "../user-alias/user-alias-dal";
|
import { TUserAliasDALFactory } from "../user-alias/user-alias-dal";
|
||||||
import { UserAliasType } from "../user-alias/user-alias-types";
|
import { UserAliasType } from "../user-alias/user-alias-types";
|
||||||
|
import { TInvalidateCacheQueueFactory } from "./invalidate-cache-queue";
|
||||||
import { TSuperAdminDALFactory } from "./super-admin-dal";
|
import { TSuperAdminDALFactory } from "./super-admin-dal";
|
||||||
import {
|
import {
|
||||||
CacheType,
|
CacheType,
|
||||||
@@ -50,6 +51,7 @@ type TSuperAdminServiceFactoryDep = {
|
|||||||
keyStore: Pick<TKeyStoreFactory, "getItem" | "setItemWithExpiry" | "deleteItem" | "deleteItems">;
|
keyStore: Pick<TKeyStoreFactory, "getItem" | "setItemWithExpiry" | "deleteItem" | "deleteItems">;
|
||||||
licenseService: Pick<TLicenseServiceFactory, "onPremFeatures">;
|
licenseService: Pick<TLicenseServiceFactory, "onPremFeatures">;
|
||||||
microsoftTeamsService: Pick<TMicrosoftTeamsServiceFactory, "initializeTeamsBot">;
|
microsoftTeamsService: Pick<TMicrosoftTeamsServiceFactory, "initializeTeamsBot">;
|
||||||
|
invalidateCacheQueue: TInvalidateCacheQueueFactory;
|
||||||
};
|
};
|
||||||
|
|
||||||
export type TSuperAdminServiceFactory = ReturnType<typeof superAdminServiceFactory>;
|
export type TSuperAdminServiceFactory = ReturnType<typeof superAdminServiceFactory>;
|
||||||
@@ -81,7 +83,8 @@ export const superAdminServiceFactory = ({
|
|||||||
identityAccessTokenDAL,
|
identityAccessTokenDAL,
|
||||||
identityTokenAuthDAL,
|
identityTokenAuthDAL,
|
||||||
identityOrgMembershipDAL,
|
identityOrgMembershipDAL,
|
||||||
microsoftTeamsService
|
microsoftTeamsService,
|
||||||
|
invalidateCacheQueue
|
||||||
}: TSuperAdminServiceFactoryDep) => {
|
}: TSuperAdminServiceFactoryDep) => {
|
||||||
const initServerCfg = async () => {
|
const initServerCfg = async () => {
|
||||||
// TODO(akhilmhdh): bad pattern time less change this later to me itself
|
// TODO(akhilmhdh): bad pattern time less change this later to me itself
|
||||||
@@ -633,12 +636,9 @@ export const superAdminServiceFactory = ({
|
|||||||
};
|
};
|
||||||
|
|
||||||
const invalidateCache = async (type: CacheType) => {
|
const invalidateCache = async (type: CacheType) => {
|
||||||
let totalKeysCleared = 0;
|
await invalidateCacheQueue.startInvalidate({
|
||||||
|
data: { type }
|
||||||
if (type === CacheType.ALL || type === CacheType.SECRETS)
|
});
|
||||||
totalKeysCleared += await keyStore.deleteItems({ pattern: "secret-manager:*" });
|
|
||||||
|
|
||||||
return totalKeysCleared;
|
|
||||||
};
|
};
|
||||||
|
|
||||||
return {
|
return {
|
||||||
|
|||||||
@@ -207,7 +207,6 @@ export type TIssueCertificateEvent = {
|
|||||||
export type TInvalidateCacheEvent = {
|
export type TInvalidateCacheEvent = {
|
||||||
event: PostHogEventTypes.InvalidateCache;
|
event: PostHogEventTypes.InvalidateCache;
|
||||||
properties: {
|
properties: {
|
||||||
keysCleared: number;
|
|
||||||
userAgent?: string;
|
userAgent?: string;
|
||||||
};
|
};
|
||||||
};
|
};
|
||||||
|
|||||||
Reference in New Issue
Block a user