diff --git a/backend/src/db/instance.ts b/backend/src/db/instance.ts index 5a8dd3d05..cdb8c3028 100644 --- a/backend/src/db/instance.ts +++ b/backend/src/db/instance.ts @@ -1,6 +1,6 @@ import knex, { Knex } from "knex"; -export type TDbClient = ReturnType; +export type TDbClient = Knex; export const initDbConnection = ({ dbConnectionUri, dbRootCert, diff --git a/backend/src/keystore/keystore.ts b/backend/src/keystore/keystore.ts index d3f19168a..6a63af776 100644 --- a/backend/src/keystore/keystore.ts +++ b/backend/src/keystore/keystore.ts @@ -2,7 +2,7 @@ import { buildRedisFromConfig, TRedisConfigKeys } from "@app/lib/config/redis"; import { pgAdvisoryLockHashText } from "@app/lib/crypto/hashtext"; import { applyJitter } from "@app/lib/dates"; import { delay as delayMs } from "@app/lib/delay"; -import { Redlock, Settings } from "@app/lib/red-lock"; +import { ExecutionResult, Redlock, Settings } from "@app/lib/red-lock"; export const PgSqlLock = { BootUpMigration: 2023, @@ -14,8 +14,6 @@ export const PgSqlLock = { CreateProject: (orgId: string) => pgAdvisoryLockHashText(`create-project:${orgId}`) } as const; -export type TKeyStoreFactory = ReturnType; - // all the key prefixes used must be set here to avoid conflict export const KeyStorePrefixes = { SecretReplication: "secret-replication-import-lock", @@ -71,7 +69,28 @@ type TWaitTillReady = { jitter?: number; }; -export const keyStoreFactory = (redisConfigKeys: TRedisConfigKeys) => { +export type TKeyStoreFactory = { + setItem: (key: string, value: string | number | Buffer, prefix?: string) => Promise<"OK">; + getItem: (key: string, prefix?: string) => Promise; + setExpiry: (key: string, expiryInSeconds: number) => Promise; + setItemWithExpiry: ( + key: string, + expiryInSeconds: number | string, + value: string | number | Buffer, + prefix?: string + ) => Promise<"OK">; + deleteItem: (key: string) => Promise; + deleteItems: (arg: TDeleteItems) => Promise; + incrementBy: (key: string, value: number) => Promise; + acquireLock( + resources: string[], + duration: number, + settings?: Partial + ): Promise<{ release: () => Promise }>; + waitTillReady: ({ key, waitingCb, keyCheckCb, waitIteration, delay, jitter }: TWaitTillReady) => Promise; +}; + +export const keyStoreFactory = (redisConfigKeys: TRedisConfigKeys): TKeyStoreFactory => { const redis = buildRedisFromConfig(redisConfigKeys); const redisLock = new Redlock([redis], { retryCount: 2, retryDelay: 200 }); @@ -108,7 +127,6 @@ export const keyStoreFactory = (redisConfigKeys: TRedisConfigKeys) => { // eslint-disable-next-line no-await-in-loop await pipeline.exec(); totalDeleted += batch.length; - console.log("BATCH DONE"); // eslint-disable-next-line no-await-in-loop await delayMs(Math.max(0, applyJitter(delay, jitter))); diff --git a/backend/src/lib/knex/index.ts b/backend/src/lib/knex/index.ts index 5949afe33..090df561a 100644 --- a/backend/src/lib/knex/index.ts +++ b/backend/src/lib/knex/index.ts @@ -71,8 +71,8 @@ export const buildFindFilter = return bd; }; -export type TFindReturn = Array< - Awaited[0] & +export type TFindReturn = Array< + Tables[Tname]["base"] & (TCount extends true ? { count: string; @@ -94,40 +94,82 @@ export type TFindOpt< tx?: Knex; }; +export type TOrmify = { + transaction: (cb: (tx: Knex) => Promise) => Promise; + findById: (id: string, tx?: Knex) => Promise; + find: ( + filter: TFindFilter, + { offset, limit, sort, count, tx, countDistinct }?: TFindOpt + ) => Promise>; + findOne: (filter: Partial, tx?: Knex) => Promise; + create: (data: Tables[Tname]["insert"], tx?: Knex) => Promise; + insertMany: (data: readonly Tables[Tname]["insert"][], tx?: Knex) => Promise; + batchInsert: (data: readonly Tables[Tname]["insert"][], tx?: Knex) => Promise; + upsert: ( + data: readonly Tables[Tname]["insert"][], + onConflictField: keyof Tables[Tname]["base"] | Array, + tx?: Knex, + mergeColumns?: (keyof Knex.ResolveTableType, "update">)[] | undefined + ) => Promise; + updateById: ( + id: string, + { + $incr, + $decr, + ...data + }: Tables[Tname]["update"] & { + $incr?: { [x in keyof Partial]: number }; + $decr?: { [x in keyof Partial]: number }; + }, + tx?: Knex + ) => Promise; + update: ( + filter: TFindFilter, + { + $incr, + $decr, + ...data + }: Tables[Tname]["update"] & { + $incr?: { [x in keyof Partial]: number }; + $decr?: { [x in keyof Partial]: number }; + }, + tx?: Knex + ) => Promise; + deleteById: (id: string, tx?: Knex) => Promise; + countDocuments: (tx?: Knex) => Promise; + delete: (filter: TFindFilter, tx?: Knex) => Promise; +}; + // What is ormify // It is to inject typical operations like find, findOne, update, delete, create // This will avoid writing most common ones each time -export const ormify = (db: Knex, tableName: Tname, dal?: DbOps) => ({ +export const ormify = ( + db: Knex, + tableName: Tname, + dal?: DbOps +): TOrmify => ({ transaction: async (cb: (tx: Knex) => Promise) => db.transaction(async (trx) => { const res = await cb(trx); return res; }), - findById: async (id: string, tx?: Knex) => { + findById: async (id, tx): Promise => { try { const result = await (tx || db.replicaNode())(tableName) .where({ id } as never) .first("*"); - return result; + return result as Tables[Tname]["base"]; } catch (error) { throw new DatabaseError({ error, name: "Find by id" }); } }, - findOne: async (filter: Partial, tx?: Knex) => { - try { - const res = await (tx || db.replicaNode())(tableName).where(filter).first("*"); - return res; - } catch (error) { - throw new DatabaseError({ error, name: "Find one" }); - } - }, find: async < TCount extends boolean = false, TCountDistinct extends keyof Tables[Tname]["base"] | undefined = undefined >( filter: TFindFilter, { offset, limit, sort, count, tx, countDistinct }: TFindOpt = {} - ) => { + ): Promise> => { try { const query = (tx || db.replicaNode())(tableName).where(buildFindFilter(filter)); if (countDistinct) { @@ -142,35 +184,43 @@ export const ormify = (db: Kne void query.orderBy(sort.map(([column, order, nulls]) => ({ column: column as string, order, nulls }))); } - const res = (await query) as TFindReturn; - return res; + const res = await query; + return res as TFindReturn; } catch (error) { throw new DatabaseError({ error, name: "Find one" }); } }, - create: async (data: Tables[Tname]["insert"], tx?: Knex) => { + findOne: async (filter, tx): Promise => { + try { + const res = await (tx || db.replicaNode())(tableName).where(filter).first("*"); + return res as Tables[Tname]["base"]; + } catch (error) { + throw new DatabaseError({ error, name: "Find one" }); + } + }, + create: async (data, tx): Promise => { try { const [res] = await (tx || db)(tableName) .insert(data as never) .returning("*"); - return res; + return res as Tables[Tname]["base"]; } catch (error) { throw new DatabaseError({ error, name: "Create" }); } }, - insertMany: async (data: readonly Tables[Tname]["insert"][], tx?: Knex) => { + insertMany: async (data, tx?): Promise => { try { if (!data.length) return []; const res = await (tx || db)(tableName) .insert(data as never) .returning("*"); - return res; + return res as Tables[Tname]["base"][]; } catch (error) { throw new DatabaseError({ error, name: "Create" }); } }, // This spilit the insert into multiple chunk - batchInsert: async (data: readonly Tables[Tname]["insert"][], tx?: Knex) => { + batchInsert: async (data, tx): Promise => { try { if (!data.length) return []; const res = await (tx || db).batchInsert(tableName, data as never).returning("*"); @@ -179,12 +229,7 @@ export const ormify = (db: Kne throw new DatabaseError({ error, name: "batchInsert" }); } }, - upsert: async ( - data: readonly Tables[Tname]["insert"][], - onConflictField: keyof Tables[Tname]["base"] | Array, - tx?: Knex, - mergeColumns?: (keyof Knex.ResolveTableType, "update">)[] | undefined - ) => { + upsert: async (data, onConflictField, tx, mergeColumns): Promise => { try { if (!data.length) return []; const res = await (tx || db)(tableName) @@ -192,23 +237,12 @@ export const ormify = (db: Kne .onConflict(onConflictField as never) .merge(mergeColumns) .returning("*"); - return res; + return res as Tables[Tname]["base"][]; } catch (error) { throw new DatabaseError({ error, name: "Create" }); } }, - updateById: async ( - id: string, - { - $incr, - $decr, - ...data - }: Tables[Tname]["update"] & { - $incr?: { [x in keyof Partial]: number }; - $decr?: { [x in keyof Partial]: number }; - }, - tx?: Knex - ) => { + updateById: async (id, { $incr, $decr, ...data }, tx): Promise => { try { const query = (tx || db)(tableName) .where({ id } as never) @@ -225,23 +259,12 @@ export const ormify = (db: Kne }); } const [docs] = await query; - return docs; + return docs as Tables[Tname]["base"]; } catch (error) { throw new DatabaseError({ error, name: "Update by id" }); } }, - update: async ( - filter: TFindFilter, - { - $incr, - $decr, - ...data - }: Tables[Tname]["update"] & { - $incr?: { [x in keyof Partial]: number }; - $decr?: { [x in keyof Partial]: number }; - }, - tx?: Knex - ) => { + update: async (filter, { $incr, $decr, ...data }, tx): Promise => { try { const query = (tx || db)(tableName) .where(buildFindFilter(filter)) @@ -258,26 +281,34 @@ export const ormify = (db: Kne void query.increment(incrementField, incrementValue); }); } - return await query; + return (await query) as Tables[Tname]["base"][]; } catch (error) { throw new DatabaseError({ error, name: "Update" }); } }, - deleteById: async (id: string, tx?: Knex) => { + deleteById: async (id, tx): Promise => { try { const [res] = await (tx || db)(tableName) .where({ id } as never) .delete() .returning("*"); - return res; + return res as Tables[Tname]["base"]; } catch (error) { throw new DatabaseError({ error, name: "Delete by id" }); } }, - delete: async (filter: TFindFilter, tx?: Knex) => { + countDocuments: async (tx): Promise => { + try { + const [res] = await (tx || db)(tableName).count({ count: "*" }).returning("*"); + return Number((res as { count: number }).count || 0); + } catch (error) { + throw new DatabaseError({ error, name: "Delete by id" }); + } + }, + delete: async (filter, tx): Promise => { try { const res = await (tx || db)(tableName).where(buildFindFilter(filter)).delete().returning("*"); - return res; + return res as Tables[Tname]["base"][]; } catch (error) { throw new DatabaseError({ error, name: "Delete" }); } diff --git a/backend/src/queue/queue-service.ts b/backend/src/queue/queue-service.ts index 25677841d..16c5bb38f 100644 --- a/backend/src/queue/queue-service.ts +++ b/backend/src/queue/queue-service.ts @@ -325,11 +325,69 @@ const isQueueEnabled = (name: QueueName) => { } }; -export type TQueueServiceFactory = ReturnType; +export type TQueueServiceFactory = { + initialize: () => Promise; + start: ( + name: T, + jobFn: (job: Job, token?: string) => Promise, + queueSettings?: Omit + ) => void; + startPg: ( + jobName: QueueJobs, + jobsFn: (jobs: PgBoss.JobWithMetadata[]) => Promise, + options: WorkOptions & { + workerCount: number; + } + ) => Promise; + listen: < + T extends QueueName, + U extends keyof WorkerListener + >( + name: T, + event: U, + listener: WorkerListener[U] + ) => void; + queue: ( + name: T, + job: TQueueJobTypes[T]["name"], + data: TQueueJobTypes[T]["payload"], + opts?: JobsOptions & { + jobId?: string; + } + ) => Promise; + queuePg: ( + job: TQueueJobTypes[T]["name"], + data: TQueueJobTypes[T]["payload"], + opts?: PgBoss.SendOptions & { jobId?: string } + ) => Promise; + schedulePg: ( + job: TQueueJobTypes[T]["name"], + cron: string, + data: TQueueJobTypes[T]["payload"], + opts?: PgBoss.ScheduleOptions & { jobId?: string } + ) => Promise; + shutdown: () => Promise; + stopRepeatableJob: ( + name: T, + job: TQueueJobTypes[T]["name"], + repeatOpt: RepeatOptions, + jobId?: string + ) => Promise; + stopRepeatableJobByJobId: (name: T, jobId: string) => Promise; + stopRepeatableJobByKey: (name: T, repeatJobKey: string) => Promise; + clearQueue: (name: QueueName) => Promise; + stopJobById: (name: T, jobId: string) => Promise; + getRepeatableJobs: ( + name: QueueName, + startOffset?: number, + endOffset?: number + ) => Promise<{ key: string; name: string; id: string | null }[]>; +}; + export const queueServiceFactory = ( redisCfg: TRedisConfigKeys, { dbConnectionUrl, dbRootCert }: { dbConnectionUrl: string; dbRootCert?: string } -) => { +): TQueueServiceFactory => { const connection = buildRedisFromConfig(redisCfg); const queueContainer = {} as Record< QueueName, @@ -366,36 +424,26 @@ export const queueServiceFactory = ( }); }; - const start = ( - name: T, - jobFn: (job: Job, token?: string) => Promise, - queueSettings: Omit = {} - ) => { + const start: TQueueServiceFactory["start"] = (name, jobFn, queueSettings) => { if (queueContainer[name]) { throw new Error(`${name} queue is already initialized`); } - queueContainer[name] = new Queue(name as string, { + queueContainer[name] = new Queue(name as string, { ...queueSettings, connection }); const appCfg = getConfig(); if (appCfg.QUEUE_WORKERS_ENABLED && isQueueEnabled(name)) { - workerContainer[name] = new Worker(name, jobFn, { + workerContainer[name] = new Worker(name, jobFn, { ...queueSettings, connection }); } }; - const startPg = async ( - jobName: QueueJobs, - jobsFn: (jobs: PgBoss.JobWithMetadata[]) => Promise, - options: WorkOptions & { - workerCount: number; - } - ) => { + const startPg: TQueueServiceFactory["startPg"] = async (jobName, jobsFn, options) => { if (queueContainerPg[jobName]) { throw new Error(`${jobName} queue is already initialized`); } @@ -429,19 +477,12 @@ export const queueServiceFactory = ( await Promise.all( Array.from({ length: options.workerCount }).map(() => - pgBoss.work(jobName, { ...options, includeMetadata: true }, jobsFn) + pgBoss.work(jobName, { ...options, includeMetadata: true }, jobsFn) ) ); }; - const listen = < - T extends QueueName, - U extends keyof WorkerListener - >( - name: T, - event: U, - listener: WorkerListener[U] - ) => { + const listen: TQueueServiceFactory["listen"] = (name, event, listener) => { const appCfg = getConfig(); if (!appCfg.QUEUE_WORKERS_ENABLED || !isQueueEnabled(name)) { return; @@ -451,12 +492,7 @@ export const queueServiceFactory = ( worker.on(event, listener); }; - const queue = async ( - name: T, - job: TQueueJobTypes[T]["name"], - data: TQueueJobTypes[T]["payload"], - opts?: JobsOptions & { jobId?: string } - ) => { + const queue: TQueueServiceFactory["queue"] = async (name, job, data, opts) => { const q = queueContainer[name]; await q.add(job, data, opts); @@ -474,35 +510,25 @@ export const queueServiceFactory = ( }); }; - const schedulePg = async ( - job: TQueueJobTypes[T]["name"], - cron: string, - data: TQueueJobTypes[T]["payload"], - opts?: PgBoss.ScheduleOptions & { jobId?: string } - ) => { + const schedulePg: TQueueServiceFactory["schedulePg"] = async (job, cron, data, opts) => { await pgBoss.schedule(job, cron, data, opts); }; - const stopRepeatableJob = async ( - name: T, - job: TQueueJobTypes[T]["name"], - repeatOpt: RepeatOptions, - jobId?: string - ) => { + const stopRepeatableJob: TQueueServiceFactory["stopRepeatableJob"] = async (name, job, repeatOpt, jobId) => { const q = queueContainer[name]; if (q) { return q.removeRepeatable(job, repeatOpt, jobId); } }; - const getRepeatableJobs = (name: QueueName, startOffset?: number, endOffset?: number) => { + const getRepeatableJobs: TQueueServiceFactory["getRepeatableJobs"] = (name, startOffset, endOffset) => { const q = queueContainer[name]; if (!q) throw new Error(`Queue '${name}' not initialized`); return q.getRepeatableJobs(startOffset, endOffset); }; - const stopRepeatableJobByJobId = async (name: T, jobId: string) => { + const stopRepeatableJobByJobId: TQueueServiceFactory["stopRepeatableJobByJobId"] = async (name, jobId) => { const q = queueContainer[name]; const job = await q.getJob(jobId); if (!job) return true; @@ -511,23 +537,23 @@ export const queueServiceFactory = ( return q.removeRepeatableByKey(job.repeatJobKey); }; - const stopRepeatableJobByKey = async (name: T, repeatJobKey: string) => { + const stopRepeatableJobByKey: TQueueServiceFactory["stopRepeatableJobByKey"] = async (name, repeatJobKey) => { const q = queueContainer[name]; return q.removeRepeatableByKey(repeatJobKey); }; - const stopJobById = async (name: T, jobId: string) => { + const stopJobById: TQueueServiceFactory["stopJobById"] = async (name, jobId) => { const q = queueContainer[name]; const job = await q.getJob(jobId); return job?.remove().catch(() => undefined); }; - const clearQueue = async (name: QueueName) => { + const clearQueue: TQueueServiceFactory["clearQueue"] = async (name) => { const q = queueContainer[name]; await q.drain(); }; - const shutdown = async () => { + const shutdown: TQueueServiceFactory["shutdown"] = async () => { await Promise.all(Object.values(workerContainer).map((worker) => worker.close())); }; diff --git a/backend/src/services/pki-templates/pki-templates-dal.ts b/backend/src/services/pki-templates/pki-templates-dal.ts index 45c632d70..4f618c807 100644 --- a/backend/src/services/pki-templates/pki-templates-dal.ts +++ b/backend/src/services/pki-templates/pki-templates-dal.ts @@ -4,7 +4,7 @@ import { Tables } from "knex/types/tables"; import { TDbClient } from "@app/db"; import { TableName } from "@app/db/schemas"; import { DatabaseError } from "@app/lib/errors"; -import { buildFindFilter, ormify, selectAllTableCols, TFindFilter, TFindOpt, TFindReturn } from "@app/lib/knex"; +import { buildFindFilter, ormify, selectAllTableCols, TFindFilter, TFindOpt } from "@app/lib/knex"; export type TPkiTemplatesDALFactory = ReturnType; @@ -91,7 +91,7 @@ export const pkiTemplatesDALFactory = (db: TDbClient) => { void query.orderBy(sort.map(([column, order, nulls]) => ({ column: column as string, order, nulls }))); } - const res = (await query) as TFindReturn; + const res = await query; return res.map((el) => ({ ...el, ca: { id: el.caId, name: el.caName } })); } catch (error) { throw new DatabaseError({ error, name: "Find one" });