feat: updated orm, keystore and queue

This commit is contained in:
=
2025-06-13 17:40:51 +05:30
parent b26e56c97e
commit 5b51ab3216
5 changed files with 190 additions and 115 deletions

View File

@@ -1,6 +1,6 @@
import knex, { Knex } from "knex";
export type TDbClient = ReturnType<typeof initDbConnection>;
export type TDbClient = Knex;
export const initDbConnection = ({
dbConnectionUri,
dbRootCert,

View File

@@ -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<typeof keyStoreFactory>;
// 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<string | null>;
setExpiry: (key: string, expiryInSeconds: number) => Promise<number>;
setItemWithExpiry: (
key: string,
expiryInSeconds: number | string,
value: string | number | Buffer,
prefix?: string
) => Promise<"OK">;
deleteItem: (key: string) => Promise<number>;
deleteItems: (arg: TDeleteItems) => Promise<number>;
incrementBy: (key: string, value: number) => Promise<number>;
acquireLock(
resources: string[],
duration: number,
settings?: Partial<Settings>
): Promise<{ release: () => Promise<ExecutionResult> }>;
waitTillReady: ({ key, waitingCb, keyCheckCb, waitIteration, delay, jitter }: TWaitTillReady) => Promise<void>;
};
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)));

View File

@@ -71,8 +71,8 @@ export const buildFindFilter =
return bd;
};
export type TFindReturn<TQuery extends Knex.QueryBuilder, TCount extends boolean = false> = Array<
Awaited<TQuery>[0] &
export type TFindReturn<Tname extends keyof Tables, TCount extends boolean = false> = Array<
Tables[Tname]["base"] &
(TCount extends true
? {
count: string;
@@ -94,40 +94,82 @@ export type TFindOpt<
tx?: Knex;
};
export type TOrmify<Tname extends keyof Tables> = {
transaction: <T>(cb: (tx: Knex) => Promise<T>) => Promise<T>;
findById: (id: string, tx?: Knex) => Promise<Tables[Tname]["base"]>;
find: <TCount extends boolean = false, TCountDistinct extends keyof Tables[Tname]["base"] | undefined = undefined>(
filter: TFindFilter<Tables[Tname]["base"]>,
{ offset, limit, sort, count, tx, countDistinct }?: TFindOpt<Tables[Tname]["base"], TCount, TCountDistinct>
) => Promise<TFindReturn<Tname, TCountDistinct extends undefined ? TCount : true>>;
findOne: (filter: Partial<Tables[Tname]["base"]>, tx?: Knex) => Promise<Tables[Tname]["base"]>;
create: (data: Tables[Tname]["insert"], tx?: Knex) => Promise<Tables[Tname]["base"]>;
insertMany: (data: readonly Tables[Tname]["insert"][], tx?: Knex) => Promise<Tables[Tname]["base"][]>;
batchInsert: (data: readonly Tables[Tname]["insert"][], tx?: Knex) => Promise<Tables[Tname]["base"][]>;
upsert: (
data: readonly Tables[Tname]["insert"][],
onConflictField: keyof Tables[Tname]["base"] | Array<keyof Tables[Tname]["base"]>,
tx?: Knex,
mergeColumns?: (keyof Knex.ResolveTableType<Knex.TableType<Tname>, "update">)[] | undefined
) => Promise<Tables[Tname]["base"][]>;
updateById: (
id: string,
{
$incr,
$decr,
...data
}: Tables[Tname]["update"] & {
$incr?: { [x in keyof Partial<Tables[Tname]["base"]>]: number };
$decr?: { [x in keyof Partial<Tables[Tname]["base"]>]: number };
},
tx?: Knex
) => Promise<Tables[Tname]["base"]>;
update: (
filter: TFindFilter<Tables[Tname]["base"]>,
{
$incr,
$decr,
...data
}: Tables[Tname]["update"] & {
$incr?: { [x in keyof Partial<Tables[Tname]["base"]>]: number };
$decr?: { [x in keyof Partial<Tables[Tname]["base"]>]: number };
},
tx?: Knex
) => Promise<Tables[Tname]["base"][]>;
deleteById: (id: string, tx?: Knex) => Promise<Tables[Tname]["base"]>;
countDocuments: (tx?: Knex) => Promise<number>;
delete: (filter: TFindFilter<Tables[Tname]["base"]>, tx?: Knex) => Promise<Tables[Tname]["base"][]>;
};
// 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 = <DbOps extends object, Tname extends keyof Tables>(db: Knex, tableName: Tname, dal?: DbOps) => ({
export const ormify = <DbOps extends object, Tname extends keyof Tables>(
db: Knex,
tableName: Tname,
dal?: DbOps
): TOrmify<Tname> => ({
transaction: async <T>(cb: (tx: Knex) => Promise<T>) =>
db.transaction(async (trx) => {
const res = await cb(trx);
return res;
}),
findById: async (id: string, tx?: Knex) => {
findById: async (id, tx): Promise<Tables[Tname]["base"]> => {
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<Tables[Tname]["base"]>, 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<Tables[Tname]["base"]>,
{ offset, limit, sort, count, tx, countDistinct }: TFindOpt<Tables[Tname]["base"], TCount, TCountDistinct> = {}
) => {
): Promise<TFindReturn<Tname, TCountDistinct extends undefined ? TCount : true>> => {
try {
const query = (tx || db.replicaNode())(tableName).where(buildFindFilter(filter));
if (countDistinct) {
@@ -142,35 +184,43 @@ export const ormify = <DbOps extends object, Tname extends keyof Tables>(db: Kne
void query.orderBy(sort.map(([column, order, nulls]) => ({ column: column as string, order, nulls })));
}
const res = (await query) as TFindReturn<typeof query, TCountDistinct extends undefined ? TCount : true>;
return res;
const res = await query;
return res as TFindReturn<Tname, TCountDistinct extends undefined ? TCount : true>;
} catch (error) {
throw new DatabaseError({ error, name: "Find one" });
}
},
create: async (data: Tables[Tname]["insert"], tx?: Knex) => {
findOne: async (filter, tx): Promise<Tables[Tname]["base"]> => {
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<Tables[Tname]["base"]> => {
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<Tables[Tname]["base"][]> => {
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<Tables[Tname]["base"][]> => {
try {
if (!data.length) return [];
const res = await (tx || db).batchInsert(tableName, data as never).returning("*");
@@ -179,12 +229,7 @@ export const ormify = <DbOps extends object, Tname extends keyof Tables>(db: Kne
throw new DatabaseError({ error, name: "batchInsert" });
}
},
upsert: async (
data: readonly Tables[Tname]["insert"][],
onConflictField: keyof Tables[Tname]["base"] | Array<keyof Tables[Tname]["base"]>,
tx?: Knex,
mergeColumns?: (keyof Knex.ResolveTableType<Knex.TableType<Tname>, "update">)[] | undefined
) => {
upsert: async (data, onConflictField, tx, mergeColumns): Promise<Tables[Tname]["base"][]> => {
try {
if (!data.length) return [];
const res = await (tx || db)(tableName)
@@ -192,23 +237,12 @@ export const ormify = <DbOps extends object, Tname extends keyof Tables>(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<Tables[Tname]["base"]>]: number };
$decr?: { [x in keyof Partial<Tables[Tname]["base"]>]: number };
},
tx?: Knex
) => {
updateById: async (id, { $incr, $decr, ...data }, tx): Promise<Tables[Tname]["base"]> => {
try {
const query = (tx || db)(tableName)
.where({ id } as never)
@@ -225,23 +259,12 @@ export const ormify = <DbOps extends object, Tname extends keyof Tables>(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<Tables[Tname]["base"]>,
{
$incr,
$decr,
...data
}: Tables[Tname]["update"] & {
$incr?: { [x in keyof Partial<Tables[Tname]["base"]>]: number };
$decr?: { [x in keyof Partial<Tables[Tname]["base"]>]: number };
},
tx?: Knex
) => {
update: async (filter, { $incr, $decr, ...data }, tx): Promise<Tables[Tname]["base"][]> => {
try {
const query = (tx || db)(tableName)
.where(buildFindFilter(filter))
@@ -258,26 +281,34 @@ export const ormify = <DbOps extends object, Tname extends keyof Tables>(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<Tables[Tname]["base"]> => {
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<Tables[Tname]["base"]>, tx?: Knex) => {
countDocuments: async (tx): Promise<number> => {
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<Tables[Tname]["base"][]> => {
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" });
}

View File

@@ -325,11 +325,69 @@ const isQueueEnabled = (name: QueueName) => {
}
};
export type TQueueServiceFactory = ReturnType<typeof queueServiceFactory>;
export type TQueueServiceFactory = {
initialize: () => Promise<void>;
start: <T extends QueueName>(
name: T,
jobFn: (job: Job<TQueueJobTypes[T]["payload"], void, TQueueJobTypes[T]["name"]>, token?: string) => Promise<void>,
queueSettings?: Omit<QueueOptions, "connection">
) => void;
startPg: <T extends QueueName>(
jobName: QueueJobs,
jobsFn: (jobs: PgBoss.JobWithMetadata<TQueueJobTypes[T]["payload"]>[]) => Promise<void>,
options: WorkOptions & {
workerCount: number;
}
) => Promise<void>;
listen: <
T extends QueueName,
U extends keyof WorkerListener<TQueueJobTypes[T]["payload"], void, TQueueJobTypes[T]["name"]>
>(
name: T,
event: U,
listener: WorkerListener<TQueueJobTypes[T]["payload"], void, TQueueJobTypes[T]["name"]>[U]
) => void;
queue: <T extends QueueName>(
name: T,
job: TQueueJobTypes[T]["name"],
data: TQueueJobTypes[T]["payload"],
opts?: JobsOptions & {
jobId?: string;
}
) => Promise<void>;
queuePg: <T extends QueueName>(
job: TQueueJobTypes[T]["name"],
data: TQueueJobTypes[T]["payload"],
opts?: PgBoss.SendOptions & { jobId?: string }
) => Promise<void>;
schedulePg: <T extends QueueName>(
job: TQueueJobTypes[T]["name"],
cron: string,
data: TQueueJobTypes[T]["payload"],
opts?: PgBoss.ScheduleOptions & { jobId?: string }
) => Promise<void>;
shutdown: () => Promise<void>;
stopRepeatableJob: <T extends QueueName>(
name: T,
job: TQueueJobTypes[T]["name"],
repeatOpt: RepeatOptions,
jobId?: string
) => Promise<boolean | undefined>;
stopRepeatableJobByJobId: <T extends QueueName>(name: T, jobId: string) => Promise<boolean>;
stopRepeatableJobByKey: <T extends QueueName>(name: T, repeatJobKey: string) => Promise<boolean>;
clearQueue: (name: QueueName) => Promise<void>;
stopJobById: <T extends QueueName>(name: T, jobId: string) => Promise<void | undefined>;
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 = <T extends QueueName>(
name: T,
jobFn: (job: Job<TQueueJobTypes[T]["payload"], void, TQueueJobTypes[T]["name"]>, token?: string) => Promise<void>,
queueSettings: Omit<QueueOptions, "connection"> = {}
) => {
const start: TQueueServiceFactory["start"] = (name, jobFn, queueSettings) => {
if (queueContainer[name]) {
throw new Error(`${name} queue is already initialized`);
}
queueContainer[name] = new Queue<TQueueJobTypes[T]["payload"], void, TQueueJobTypes[T]["name"]>(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<TQueueJobTypes[T]["payload"], void, TQueueJobTypes[T]["name"]>(name, jobFn, {
workerContainer[name] = new Worker(name, jobFn, {
...queueSettings,
connection
});
}
};
const startPg = async <T extends QueueName>(
jobName: QueueJobs,
jobsFn: (jobs: PgBoss.JobWithMetadata<TQueueJobTypes[T]["payload"]>[]) => Promise<void>,
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<TQueueJobTypes[T]["payload"]>(jobName, { ...options, includeMetadata: true }, jobsFn)
pgBoss.work(jobName, { ...options, includeMetadata: true }, jobsFn)
)
);
};
const listen = <
T extends QueueName,
U extends keyof WorkerListener<TQueueJobTypes[T]["payload"], void, TQueueJobTypes[T]["name"]>
>(
name: T,
event: U,
listener: WorkerListener<TQueueJobTypes[T]["payload"], void, TQueueJobTypes[T]["name"]>[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 <T extends QueueName>(
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 <T extends QueueName>(
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 <T extends QueueName>(
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 <T extends QueueName>(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 <T extends QueueName>(name: T, repeatJobKey: string) => {
const stopRepeatableJobByKey: TQueueServiceFactory["stopRepeatableJobByKey"] = async (name, repeatJobKey) => {
const q = queueContainer[name];
return q.removeRepeatableByKey(repeatJobKey);
};
const stopJobById = async <T extends QueueName>(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()));
};

View File

@@ -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<typeof pkiTemplatesDALFactory>;
@@ -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<typeof query, TCountDistinct extends undefined ? TCount : true>;
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" });