{"version":3,"sources":["../../../../../src/dbs/adapters/mongodb/index.ts","../../../../../src/dbs/adapters/mongodb/changes.ts","../../../../../src/errors/equippedError.ts","../../../../../src/events/adapters/base.ts","../../../../../src/instance/index.ts","../../../../../src/instance/hooks.ts","../../../../../src/instance/settings.ts","../../../../../src/utilities/configurable.ts","../../../../../src/utilities/json.ts","../../../../../src/utilities/random.ts","../../../../../src/utilities/retry.ts","../../../../../src/dbs/adapters/base/changes.ts","../../../../../src/dbs/adapters/base/types.ts","../../../../../src/events/adapters/kafka/index.ts","../../../../../src/dbs/adapters/mongodb/db.ts","../../../../../src/dbs/adapters/mongodb/pipes.ts","../../../../../src/dbs/adapters/mongodb/query.ts","../../../../../src/dbs/pipes.ts","../../../../../src/dbs/adapters/base/db.ts"],"sourcesContent":["export * from './changes'\nexport * from './db'\nexport * from './pipes'\n","import { Collection, type Filter } from 'mongodb'\nimport { differ } from 'valleyed'\n\nimport type { MongoDbConfig } from './pipes'\nimport { EquippedError } from '../../../errors'\nimport { Instance } from '../../../instance'\nimport { retry } from '../../../utilities'\nimport { DbChange, TopicPrefix } from '../base/changes'\nimport * as core from '../base/core'\nimport type { DbChangeConfig } from '../base/types'\n\nexport class MongoDbChange<Model extends core.Model<{ _id: string }>, Entity extends core.Entity> extends DbChange<Model, Entity> {\n\t#started = false\n\n\tconstructor(\n\t\tconfig: MongoDbConfig,\n\t\tchange: DbChangeConfig,\n\t\tcollection: Collection<Model>,\n\t\tcallbacks: core.DbChangeCallbacks<Model, Entity>,\n\t\tmapper: (model: Model) => Entity,\n\t) {\n\t\tsuper(change, callbacks, mapper)\n\n\t\tconst hydrate = (data: any) =>\n\t\t\tdata._id\n\t\t\t\t? {\n\t\t\t\t\t\t...data,\n\t\t\t\t\t\t_id: data._id['$oid'] ?? data._id,\n\t\t\t\t\t}\n\t\t\t\t: undefined\n\n\t\tconst dbName = collection.dbName\n\t\tconst colName = collection.collectionName\n\t\tconst dbColName = `${dbName}.${colName}`\n\t\tconst topic = `${TopicPrefix}.${dbColName}`\n\n\t\tconst TestId = '5f5f65717569707065645f5f' // __equipped__\n\t\tconst condition = { _id: TestId } as Filter<Model>\n\n\t\tchange.eventBus.stream(topic as never, { skipScope: true }).subscribe(async (data: DbDocumentChange) => {\n\t\t\tconst op = data.op\n\n\t\t\tlet before = JSON.parse(data.before ?? 'null')\n\t\t\tlet after = JSON.parse(data.after ?? 'null')\n\n\t\t\tif (before) before = hydrate(before)\n\t\t\tif (after) after = hydrate(after)\n\t\t\tif (before?._id === TestId || after?._id === TestId) return\n\n\t\t\tif (op === 'c' && this.callbacks.created && after)\n\t\t\t\tawait this.callbacks.created({\n\t\t\t\t\tbefore: null,\n\t\t\t\t\tafter: this.mapper(after)!,\n\t\t\t\t})\n\t\t\telse if (op === 'u' && this.callbacks.updated && before && after)\n\t\t\t\tawait this.callbacks.updated({\n\t\t\t\t\tbefore: this.mapper(before)!,\n\t\t\t\t\tafter: this.mapper(after)!,\n\t\t\t\t\tchanges: differ.from(differ.diff(before, after)),\n\t\t\t\t})\n\t\t\telse if (op === 'd' && this.callbacks.deleted && before)\n\t\t\t\tawait this.callbacks.deleted({\n\t\t\t\t\tbefore: this.mapper(before)!,\n\t\t\t\t\tafter: null,\n\t\t\t\t})\n\t\t})\n\n\t\tInstance.on('start', async () => {\n\t\t\tif (this.#started) return\n\t\t\tthis.#started = true\n\n\t\t\tawait retry(\n\t\t\t\tasync () => {\n\t\t\t\t\tconst started = await this.configureConnector(topic, {\n\t\t\t\t\t\t'connector.class': 'io.debezium.connector.mongodb.MongoDbConnector',\n\t\t\t\t\t\t'capture.mode': 'change_streams_update_full_with_pre_image',\n\t\t\t\t\t\t'mongodb.connection.string': config.uri,\n\t\t\t\t\t\t'collection.include.list': dbColName,\n\t\t\t\t\t\t'snapshot.mode': 'when_needed',\n\t\t\t\t\t\t'capture.scope': 'collection',\n\t\t\t\t\t\t'capture.target': dbColName,\n\t\t\t\t\t\t'heartbeat.interval.ms': '60000',\n\t\t\t\t\t\t'errors.max.retries': '3',\n\t\t\t\t\t})\n\n\t\t\t\t\tif (started) return { done: true, value: true }\n\t\t\t\t\tawait collection.findOneAndUpdate(condition, { $set: { colName } as any }, { upsert: true })\n\t\t\t\t\tawait collection.findOneAndDelete(condition)\n\t\t\t\t\tInstance.get().log.warn(`Waiting for db changes for ${dbColName} to start...`)\n\t\t\t\t\treturn { done: false }\n\t\t\t\t},\n\t\t\t\t6,\n\t\t\t\t30_000,\n\t\t\t).catch((err) => Instance.crash(new EquippedError(`Failed to start db changes`, { dbColName }, err)))\n\t\t})\n\t}\n}\n\ntype DbDocumentChange = {\n\tbefore: string | null\n\tafter: string | null\n\top: 'c' | 'u' | 'd'\n}\n","export class EquippedError extends Error {\n\tconstructor(\n\t\tpublic readonly message: string,\n\t\tpublic readonly context: Record<string, unknown>,\n\t\tpublic readonly cause?: unknown,\n\t) {\n\t\tsuper(message, { cause })\n\t}\n}\n","import type { Events } from '../../types'\n\nexport type StreamOptions = { skipScope?: boolean; fanout: boolean }\n\nexport type Stream<EventData> = {\n\tpublish: (data: EventData) => Promise<boolean>\n\tsubscribe: (onMessage: (data: EventData) => Promise<void>) => void\n}\n\nexport abstract class EventBus {\n\tabstract stream<Event extends Events[keyof Events]>(topic: Event['topic'], options?: Partial<StreamOptions>): Stream<Event['data']>\n}\n","import pino, { type Logger } from 'pino'\nimport { ulid } from 'ulid'\nimport { type Pipe, v } from 'valleyed'\n\nimport { EquippedError } from '../errors'\nimport { type ClassRef, type HookCb, type HookEvent, type HookOptions, type HookRecord, registerHook, runHooks } from './hooks'\nimport { instanceSettingsPipe, type Settings, type SettingsInput } from './settings'\n\nexport type { ClassRef, HookCb, HookEvent, HookOptions }\n\nexport class Instance {\n\tstatic #id: string | undefined\n\tstatic #instance: Instance | undefined\n\tstatic #hooks: Partial<Record<HookEvent, HookRecord[]>> = {}\n\treadonly settings: Readonly<Settings>\n\treadonly log: Logger<never>\n\n\tprivate constructor(settings: Settings) {\n\t\tInstance.#instance = this\n\t\tthis.settings = Object.freeze(settings)\n\t\tthis.log = pino<never>({\n\t\t\tlevel: this.settings.log.level,\n\t\t\tserializers: {\n\t\t\t\terr: pino.stdSerializers.err,\n\t\t\t\terror: pino.stdSerializers.err,\n\t\t\t\treq: pino.stdSerializers.req,\n\t\t\t\tres: pino.stdSerializers.res,\n\t\t\t},\n\t\t\tmixin: () => ({\n\t\t\t\tinstanceId: Instance.#id,\n\t\t\t}),\n\t\t})\n\t\tInstance.#registerOnExitHandler()\n\t}\n\n\talias(id: string) {\n\t\tif (Instance.#id !== undefined) return Instance.crash(new EquippedError('Instance already has an alias', {}))\n\t\tInstance.#id = id\n\t}\n\n\tget id() {\n\t\tif (Instance.#id === undefined) return Instance.crash(new EquippedError('Instance doesnt have an alias yet', {}))\n\t\treturn Instance.#id\n\t}\n\n\tgetScopedName(name: string, key = '.') {\n\t\treturn [this.settings.app.name, name].join(key)\n\t}\n\n\tasync start() {\n\t\ttry {\n\t\t\tawait runHooks(Instance.#hooks['setup'] ?? [])\n\t\t\tawait runHooks(Instance.#hooks['start'] ?? [])\n\t\t} catch (error) {\n\t\t\tInstance.crash(new EquippedError(`Error starting instance`, {}, error))\n\t\t}\n\t}\n\n\tstatic envs<E extends object>(envsPipe: Pipe<unknown, E>): E {\n\t\tconst envValidity = v.validate(envsPipe, process.env)\n\t\tif (!envValidity.valid) {\n\t\t\tInstance.crash(\n\t\t\t\tnew EquippedError(`Environment variables are not valid\\n${envValidity.error.toString()}`, {\n\t\t\t\t\tmessages: envValidity.error.messages,\n\t\t\t\t}),\n\t\t\t)\n\t\t}\n\t\treturn envValidity.value\n\t}\n\n\tstatic create(settings: SettingsInput) {\n\t\tif (Instance.#instance) return Instance.crash(new EquippedError('Instance has been initialized already', {}))\n\t\tconst settingsValidity = v.validate(instanceSettingsPipe(), settings)\n\t\tif (!settingsValidity.valid) {\n\t\t\tInstance.crash(\n\t\t\t\tnew EquippedError(`Settings are not valid\\n${settingsValidity.error.toString()}`, {\n\t\t\t\t\tmessages: settingsValidity.error.messages,\n\t\t\t\t}),\n\t\t\t)\n\t\t}\n\t\treturn new Instance(settingsValidity.value)\n\t}\n\n\tstatic get() {\n\t\tif (!Instance.#instance)\n\t\t\treturn Instance.crash(\n\t\t\t\tnew EquippedError('Has not been initialized. Make sure an instance has been created before you get an instance', {}),\n\t\t\t)\n\t\treturn Instance.#instance\n\t}\n\n\tstatic maybeGet() {\n\t\treturn Instance.#instance\n\t}\n\n\tstatic on(event: HookEvent, cb: HookCb, options?: HookOptions) {\n\t\tInstance.#hooks[event] ??= []\n\t\tconst record: HookRecord = { cb, class: options?.class, after: options?.after ?? [] }\n\t\tregisterHook(Instance.#hooks[event], record)\n\t}\n\n\tstatic #registerOnExitHandler() {\n\t\tconst signals = {\n\t\t\tSIGHUP: 1,\n\t\t\tSIGINT: 2,\n\t\t\tSIGTERM: 15,\n\t\t}\n\n\t\tObject.entries(signals).forEach(([signal, code]) => {\n\t\t\tprocess.on(signal, async () => {\n\t\t\t\tawait runHooks(Instance.#hooks['close'] ?? [], () => {}, true)\n\t\t\t\tprocess.exit(128 + code)\n\t\t\t})\n\t\t})\n\t}\n\n\tstatic resolveBeforeCrash<T>(cb: () => Promise<T>) {\n\t\tconst value = cb()\n\t\tInstance.on('close', async () => await value)\n\t\treturn value\n\t}\n\n\tstatic crash(error: EquippedError): never {\n\t\t// eslint-disable-next-line no-console\n\t\tconsole.error(error)\n\t\tprocess.exit(1)\n\t}\n\n\tstatic createId(opts?: { prefix?: string; time?: Date }) {\n\t\treturn `${opts?.prefix ?? ''}${ulid(opts?.time?.getTime())}`\n\t}\n}\n","export type HookEvent = 'setup' | 'start' | 'close'\nexport type HookCb = Promise<unknown | void> | (() => void | unknown | Promise<void | unknown>)\n\nexport type ClassRef = Function & { prototype: unknown; name: string }\nexport type HookOptions = {\n\tclass?: ClassRef\n\tafter?: ClassRef[]\n}\nexport type HookRecord = { cb: HookCb; class?: ClassRef; after: ClassRef[] }\n\nexport function registerHook(hooks: HookRecord[], record: HookRecord): void {\n\tfor (const dep of record.after) {\n\t\tif (!hooks.some((h) => h.class === dep)) {\n\t\t\tconst depName = dep.name || 'unknown'\n\t\t\tconst ownerName = record.class?.name || 'anonymous'\n\t\t\tthrow new Error(`Missing dependency: ${ownerName} declares after: [${depName}], but ${depName} is not registered`)\n\t\t}\n\t}\n\n\tif (record.class) {\n\t\tconst tentative = [...hooks, record]\n\t\tdetectCycle(tentative)\n\t}\n\n\thooks.push(record)\n}\n\nfunction detectCycle(hooks: HookRecord[]): void {\n\tconst classToAfter = new Map<ClassRef, ClassRef[]>()\n\tfor (const h of hooks) {\n\t\tif (!h.class) continue\n\t\tconst existing = classToAfter.get(h.class)\n\t\tif (existing) {\n\t\t\tfor (const dep of h.after) {\n\t\t\t\tif (!existing.includes(dep)) existing.push(dep)\n\t\t\t}\n\t\t} else {\n\t\t\tclassToAfter.set(h.class, [...h.after])\n\t\t}\n\t}\n\n\tconst visited = new Set<ClassRef>()\n\tconst stack = new Set<ClassRef>()\n\n\tfunction visit(cls: ClassRef, path: ClassRef[]): void {\n\t\tif (stack.has(cls)) {\n\t\t\tconst cycleStart = path.indexOf(cls)\n\t\t\tconst cycle = [...path.slice(cycleStart), cls]\n\t\t\tconst names = cycle.map((c) => c.name || 'unknown')\n\t\t\tthrow new Error(`Cycle detected: ${names.join(' → ')}`)\n\t\t}\n\t\tif (visited.has(cls)) return\n\t\tstack.add(cls)\n\t\tpath.push(cls)\n\t\tfor (const dep of classToAfter.get(cls) ?? []) {\n\t\t\tvisit(dep, path)\n\t\t}\n\t\tpath.pop()\n\t\tstack.delete(cls)\n\t\tvisited.add(cls)\n\t}\n\n\tfor (const cls of classToAfter.keys()) {\n\t\tvisit(cls, [])\n\t}\n}\n\nexport type HookLayer = HookRecord[]\n\nexport function resolveHookDAG(hooks: HookRecord[], invert: boolean = false): HookLayer[] {\n\tif (hooks.length === 0) return []\n\n\tconst classToHooks = new Map<ClassRef, HookRecord[]>()\n\tconst anonymous: HookRecord[] = []\n\n\tfor (const h of hooks) {\n\t\tif (h.class) {\n\t\t\tconst list = classToHooks.get(h.class)\n\t\t\tif (list) list.push(h)\n\t\t\telse classToHooks.set(h.class, [h])\n\t\t} else {\n\t\t\tanonymous.push(h)\n\t\t}\n\t}\n\n\tconst classes = [...classToHooks.keys()]\n\n\tconst inDegree = new Map<ClassRef, number>()\n\tconst graph = new Map<ClassRef, ClassRef[]>()\n\n\tfor (const cls of classes) {\n\t\tinDegree.set(cls, 0)\n\t\tgraph.set(cls, [])\n\t}\n\n\tfor (const h of hooks) {\n\t\tif (!h.class) continue\n\t\tfor (const dep of h.after) {\n\t\t\tconst [from, to] = invert ? [h.class, dep] : [dep, h.class]\n\t\t\tconst edges = graph.get(from)!\n\t\t\tif (!edges.includes(to)) {\n\t\t\t\tedges.push(to)\n\t\t\t\tinDegree.set(to, inDegree.get(to)! + 1)\n\t\t\t}\n\t\t}\n\t}\n\n\tconst layers: HookLayer[] = []\n\tlet queue = classes.filter((c) => inDegree.get(c) === 0)\n\n\twhile (queue.length > 0) {\n\t\tconst layer: HookRecord[] = []\n\t\tconst next: ClassRef[] = []\n\n\t\tfor (const cls of queue) {\n\t\t\tlayer.push(...(classToHooks.get(cls) ?? []))\n\t\t}\n\t\tlayers.push(layer)\n\n\t\tfor (const cls of queue) {\n\t\t\tfor (const neighbor of graph.get(cls) ?? []) {\n\t\t\t\tconst d = inDegree.get(neighbor)! - 1\n\t\t\t\tinDegree.set(neighbor, d)\n\t\t\t\tif (d === 0) next.push(neighbor)\n\t\t\t}\n\t\t}\n\n\t\tqueue = next\n\t}\n\n\tif (anonymous.length > 0) {\n\t\tconst anonWithDeps = anonymous.filter((h) => h.after.length > 0)\n\t\tconst anonWithoutDeps = anonymous.filter((h) => h.after.length === 0)\n\n\t\tif (anonWithDeps.length > 0) {\n\t\t\tlet maxLayer = -1\n\t\t\tfor (const h of anonWithDeps) {\n\t\t\t\tfor (const dep of h.after) {\n\t\t\t\t\tconst depLayer = layers.findIndex((l) => l.some((r) => r.class === dep))\n\t\t\t\t\tif (depLayer > maxLayer) maxLayer = depLayer\n\t\t\t\t}\n\t\t\t}\n\t\t\tconst insertAt = maxLayer + 1\n\t\t\tif (insertAt < layers.length) {\n\t\t\t\tlayers[insertAt].push(...anonWithDeps)\n\t\t\t} else {\n\t\t\t\tlayers.push(anonWithDeps)\n\t\t\t}\n\t\t}\n\n\t\tif (anonWithoutDeps.length > 0) {\n\t\t\tif (layers.length > 0) {\n\t\t\t\tlayers[layers.length - 1].push(...anonWithoutDeps)\n\t\t\t} else {\n\t\t\t\tlayers.push(anonWithoutDeps)\n\t\t\t}\n\t\t}\n\t}\n\n\treturn layers\n}\n\nexport async function runHooks(\n\thooks: HookRecord[],\n\tonError: (error: Error) => void = (error) => {\n\t\tthrow error\n\t},\n\tinvert: boolean = false,\n) {\n\tconst layers = resolveHookDAG(hooks, invert)\n\tfor (const layer of layers)\n\t\tawait Promise.all(\n\t\t\tlayer.map(async (h) => {\n\t\t\t\ttry {\n\t\t\t\t\tif (typeof h.cb === 'function') return await h.cb()\n\t\t\t\t\treturn await h.cb\n\t\t\t\t} catch (error) {\n\t\t\t\t\treturn onError(error instanceof Error ? error : new Error(`${error}`))\n\t\t\t\t}\n\t\t\t}),\n\t\t)\n}\n\nif (import.meta.vitest) {\n\tconst { describe, test, expect } = import.meta.vitest\n\n\tclass A {}\n\tclass B {}\n\tclass C {}\n\tclass D {}\n\n\tconst noop = () => {}\n\n\tdescribe('registerHook', () => {\n\t\ttest('registers a hook with no dependencies', () => {\n\t\t\tconst hooks: HookRecord[] = []\n\t\t\tregisterHook(hooks, { cb: noop, class: A, after: [] })\n\t\t\texpect(hooks).toHaveLength(1)\n\t\t})\n\n\t\ttest('throws on missing dependency', () => {\n\t\t\tconst hooks: HookRecord[] = []\n\t\t\texpect(() => registerHook(hooks, { cb: noop, class: B, after: [A] })).toThrow(/Missing dependency.*B.*A/)\n\t\t})\n\n\t\ttest('throws on cycle: A→B→A', () => {\n\t\t\tconst hooks: HookRecord[] = []\n\t\t\tregisterHook(hooks, { cb: noop, class: A, after: [] })\n\t\t\tregisterHook(hooks, { cb: noop, class: B, after: [A] })\n\t\t\texpect(() => registerHook(hooks, { cb: noop, class: A, after: [B] })).toThrow(/Cycle detected/)\n\t\t})\n\n\t\ttest('throws on transitive cycle: A→B→C→A', () => {\n\t\t\tconst hooks: HookRecord[] = []\n\t\t\tregisterHook(hooks, { cb: noop, class: A, after: [] })\n\t\t\tregisterHook(hooks, { cb: noop, class: B, after: [A] })\n\t\t\tregisterHook(hooks, { cb: noop, class: C, after: [B] })\n\t\t\texpect(() => registerHook(hooks, { cb: noop, class: A, after: [C] })).toThrow(/Cycle detected/)\n\t\t})\n\n\t\ttest('allows anonymous hooks with dependencies', () => {\n\t\t\tconst hooks: HookRecord[] = []\n\t\t\tregisterHook(hooks, { cb: noop, class: A, after: [] })\n\t\t\tregisterHook(hooks, { cb: noop, after: [A] })\n\t\t\texpect(hooks).toHaveLength(2)\n\t\t})\n\n\t\ttest('anonymous hook throws on missing dependency', () => {\n\t\t\tconst hooks: HookRecord[] = []\n\t\t\texpect(() => registerHook(hooks, { cb: noop, after: [A] })).toThrow(/Missing dependency.*anonymous.*A/)\n\t\t})\n\t})\n\n\tdescribe('resolveHookDAG', () => {\n\t\ttest('linear chain: A → B → C', () => {\n\t\t\tconst hooks: HookRecord[] = [\n\t\t\t\t{ cb: noop, class: A, after: [] },\n\t\t\t\t{ cb: noop, class: B, after: [A] },\n\t\t\t\t{ cb: noop, class: C, after: [B] },\n\t\t\t]\n\t\t\tconst layers = resolveHookDAG(hooks)\n\t\t\texpect(layers).toHaveLength(3)\n\t\t\texpect(layers[0].every((h) => h.class === A)).toBe(true)\n\t\t\texpect(layers[1].every((h) => h.class === B)).toBe(true)\n\t\t\texpect(layers[2].every((h) => h.class === C)).toBe(true)\n\t\t})\n\n\t\ttest('fan-out: A → B, A → C (B and C are parallel)', () => {\n\t\t\tconst hooks: HookRecord[] = [\n\t\t\t\t{ cb: noop, class: A, after: [] },\n\t\t\t\t{ cb: noop, class: B, after: [A] },\n\t\t\t\t{ cb: noop, class: C, after: [A] },\n\t\t\t]\n\t\t\tconst layers = resolveHookDAG(hooks)\n\t\t\texpect(layers).toHaveLength(2)\n\t\t\texpect(layers[0].every((h) => h.class === A)).toBe(true)\n\t\t\tconst secondClasses = layers[1].map((h) => h.class)\n\t\t\texpect(secondClasses).toContain(B)\n\t\t\texpect(secondClasses).toContain(C)\n\t\t})\n\n\t\ttest('fan-in: B → D, C → D', () => {\n\t\t\tconst hooks: HookRecord[] = [\n\t\t\t\t{ cb: noop, class: B, after: [] },\n\t\t\t\t{ cb: noop, class: C, after: [] },\n\t\t\t\t{ cb: noop, class: D, after: [B, C] },\n\t\t\t]\n\t\t\tconst layers = resolveHookDAG(hooks)\n\t\t\texpect(layers).toHaveLength(2)\n\t\t\tconst firstClasses = layers[0].map((h) => h.class)\n\t\t\texpect(firstClasses).toContain(B)\n\t\t\texpect(firstClasses).toContain(C)\n\t\t\texpect(layers[1].every((h) => h.class === D)).toBe(true)\n\t\t})\n\n\t\ttest('diamond: A → B, A → C, B → D, C → D', () => {\n\t\t\tconst hooks: HookRecord[] = [\n\t\t\t\t{ cb: noop, class: A, after: [] },\n\t\t\t\t{ cb: noop, class: B, after: [A] },\n\t\t\t\t{ cb: noop, class: C, after: [A] },\n\t\t\t\t{ cb: noop, class: D, after: [B, C] },\n\t\t\t]\n\t\t\tconst layers = resolveHookDAG(hooks)\n\t\t\texpect(layers).toHaveLength(3)\n\t\t\texpect(layers[0][0].class).toBe(A)\n\t\t\tconst midClasses = layers[1].map((h) => h.class)\n\t\t\texpect(midClasses).toContain(B)\n\t\t\texpect(midClasses).toContain(C)\n\t\t\texpect(layers[2][0].class).toBe(D)\n\t\t})\n\n\t\ttest('close-event inversion reverses the graph', () => {\n\t\t\tconst hooks: HookRecord[] = [\n\t\t\t\t{ cb: noop, class: A, after: [] },\n\t\t\t\t{ cb: noop, class: B, after: [A] },\n\t\t\t\t{ cb: noop, class: C, after: [B] },\n\t\t\t]\n\t\t\tconst layers = resolveHookDAG(hooks, true)\n\t\t\texpect(layers).toHaveLength(3)\n\t\t\texpect(layers[0].every((h) => h.class === C)).toBe(true)\n\t\t\texpect(layers[1].every((h) => h.class === B)).toBe(true)\n\t\t\texpect(layers[2].every((h) => h.class === A)).toBe(true)\n\t\t})\n\n\t\ttest('close inversion: fan-out becomes fan-in', () => {\n\t\t\tconst hooks: HookRecord[] = [\n\t\t\t\t{ cb: noop, class: A, after: [] },\n\t\t\t\t{ cb: noop, class: B, after: [A] },\n\t\t\t\t{ cb: noop, class: C, after: [A] },\n\t\t\t]\n\t\t\tconst layers = resolveHookDAG(hooks, true)\n\t\t\texpect(layers).toHaveLength(2)\n\t\t\tconst firstClasses = layers[0].map((h) => h.class)\n\t\t\texpect(firstClasses).toContain(B)\n\t\t\texpect(firstClasses).toContain(C)\n\t\t\texpect(layers[1].every((h) => h.class === A)).toBe(true)\n\t\t})\n\n\t\ttest('multiple hooks per class are all in the same layer', () => {\n\t\t\tconst cb1 = () => {}\n\t\t\tconst cb2 = () => {}\n\t\t\tconst hooks: HookRecord[] = [\n\t\t\t\t{ cb: cb1, class: A, after: [] },\n\t\t\t\t{ cb: cb2, class: A, after: [] },\n\t\t\t\t{ cb: noop, class: B, after: [A] },\n\t\t\t]\n\t\t\tconst layers = resolveHookDAG(hooks)\n\t\t\texpect(layers).toHaveLength(2)\n\t\t\texpect(layers[0]).toHaveLength(2)\n\t\t\texpect(layers[0].every((h) => h.class === A)).toBe(true)\n\t\t})\n\n\t\ttest('anonymous hooks without deps run at deepest level', () => {\n\t\t\tconst hooks: HookRecord[] = [\n\t\t\t\t{ cb: noop, class: A, after: [] },\n\t\t\t\t{ cb: noop, class: B, after: [A] },\n\t\t\t\t{ cb: noop, after: [] },\n\t\t\t]\n\t\t\tconst layers = resolveHookDAG(hooks)\n\t\t\texpect(layers).toHaveLength(2)\n\t\t\tconst lastLayer = layers[layers.length - 1]\n\t\t\texpect(lastLayer.some((h) => h.class === undefined)).toBe(true)\n\t\t})\n\n\t\ttest('anonymous hooks with deps run after their dependencies', () => {\n\t\t\tconst hooks: HookRecord[] = [\n\t\t\t\t{ cb: noop, class: A, after: [] },\n\t\t\t\t{ cb: noop, class: B, after: [A] },\n\t\t\t\t{ cb: noop, class: C, after: [B] },\n\t\t\t\t{ cb: noop, after: [A] },\n\t\t\t]\n\t\t\tconst layers = resolveHookDAG(hooks)\n\t\t\tconst anonLayer = layers.findIndex((l) => l.some((h) => h.class === undefined))\n\t\t\tconst aLayer = layers.findIndex((l) => l.some((h) => h.class === A))\n\t\t\texpect(anonLayer).toBeGreaterThan(aLayer)\n\t\t})\n\n\t\ttest('empty hooks returns empty layers', () => {\n\t\t\texpect(resolveHookDAG([])).toEqual([])\n\t\t})\n\t})\n\n\tdescribe('runHooks', () => {\n\t\ttest('same-depth siblings observably overlap', async () => {\n\t\t\tconst log: string[] = []\n\t\t\tconst hooks: HookRecord[] = [\n\t\t\t\t{\n\t\t\t\t\tcb: async () => {\n\t\t\t\t\t\tlog.push('B-start')\n\t\t\t\t\t\tawait new Promise((r) => setTimeout(r, 20))\n\t\t\t\t\t\tlog.push('B-end')\n\t\t\t\t\t},\n\t\t\t\t\tclass: B,\n\t\t\t\t\tafter: [],\n\t\t\t\t},\n\t\t\t\t{\n\t\t\t\t\tcb: async () => {\n\t\t\t\t\t\tlog.push('C-start')\n\t\t\t\t\t\tawait new Promise((r) => setTimeout(r, 20))\n\t\t\t\t\t\tlog.push('C-end')\n\t\t\t\t\t},\n\t\t\t\t\tclass: C,\n\t\t\t\t\tafter: [],\n\t\t\t\t},\n\t\t\t]\n\t\t\tawait runHooks(hooks)\n\t\t\texpect(log[0]).toBe('B-start')\n\t\t\texpect(log[1]).toBe('C-start')\n\t\t})\n\n\t\ttest('different depths run sequentially', async () => {\n\t\t\tconst log: string[] = []\n\t\t\tconst hooks: HookRecord[] = [\n\t\t\t\t{\n\t\t\t\t\tcb: async () => {\n\t\t\t\t\t\tlog.push('A')\n\t\t\t\t\t},\n\t\t\t\t\tclass: A,\n\t\t\t\t\tafter: [],\n\t\t\t\t},\n\t\t\t\t{\n\t\t\t\t\tcb: async () => {\n\t\t\t\t\t\tlog.push('B')\n\t\t\t\t\t},\n\t\t\t\t\tclass: B,\n\t\t\t\t\tafter: [A],\n\t\t\t\t},\n\t\t\t]\n\t\t\tawait runHooks(hooks)\n\t\t\texpect(log).toEqual(['A', 'B'])\n\t\t})\n\n\t\ttest('close inversion runs hooks in reverse dependency order', async () => {\n\t\t\tconst log: string[] = []\n\t\t\tconst hooks: HookRecord[] = [\n\t\t\t\t{ cb: async () => log.push('A'), class: A, after: [] },\n\t\t\t\t{ cb: async () => log.push('B'), class: B, after: [A] },\n\t\t\t\t{ cb: async () => log.push('C'), class: C, after: [B] },\n\t\t\t]\n\t\t\tawait runHooks(hooks, undefined, true)\n\t\t\texpect(log).toEqual(['C', 'B', 'A'])\n\t\t})\n\n\t\ttest('errors are routed to onError handler', async () => {\n\t\t\tconst errors: Error[] = []\n\t\t\tconst hooks: HookRecord[] = [\n\t\t\t\t{\n\t\t\t\t\tcb: () => {\n\t\t\t\t\t\tthrow new Error('boom')\n\t\t\t\t\t},\n\t\t\t\t\tclass: A,\n\t\t\t\t\tafter: [],\n\t\t\t\t},\n\t\t\t]\n\t\t\tawait runHooks(hooks, (e) => {\n\t\t\t\terrors.push(e)\n\t\t\t})\n\t\t\texpect(errors).toHaveLength(1)\n\t\t\texpect(errors[0].message).toBe('boom')\n\t\t})\n\n\t\ttest('raw promise callbacks are awaited', async () => {\n\t\t\tconst log: string[] = []\n\t\t\tconst hooks: HookRecord[] = [\n\t\t\t\t{\n\t\t\t\t\tcb: Promise.resolve().then(() => log.push('resolved')),\n\t\t\t\t\tclass: A,\n\t\t\t\t\tafter: [],\n\t\t\t\t},\n\t\t\t]\n\t\t\tawait runHooks(hooks)\n\t\t\texpect(log).toContain('resolved')\n\t\t})\n\t})\n}\n","import { type ConditionalObjectKeys, type PipeInput, type PipeOutput, v } from 'valleyed'\n\nexport const instanceSettingsPipe = () =>\n\tv.object({\n\t\tapp: v.object({\n\t\t\tname: v.string(),\n\t\t}),\n\t\tlog: v.defaults(\n\t\t\tv.object({\n\t\t\t\tlevel: v.defaults(v.in(['fatal', 'error', 'warn', 'info', 'debug', 'trace', 'silent'] as const), 'info'),\n\t\t\t}),\n\t\t\t{},\n\t\t),\n\t\tutils: v.defaults(\n\t\t\tv.object({\n\t\t\t\thashSaltRounds: v.defaults(v.number(), 10),\n\t\t\t\tpaginationDefaultLimit: v.defaults(v.number(), 100),\n\t\t\t\tmaxFileUploadSizeInMb: v.defaults(v.number(), 500),\n\t\t\t}),\n\t\t\t{},\n\t\t),\n\t})\n\nexport type Settings = PipeOutput<ReturnType<typeof instanceSettingsPipe>>\nexport type SettingsInput = ConditionalObjectKeys<PipeInput<ReturnType<typeof instanceSettingsPipe>>>\n","import { v, type ConditionalObjectKeys, type Pipe, type PipeInput, type PipeOutput } from 'valleyed'\n\ntype CtorParams<T> = ConstructorParameters<T & (abstract new (...args: any[]) => any)>\ntype BaseCtorParams<T> = T extends abstract new (...args: infer A) => any ? A : never\n\nexport function configurable<P extends Pipe<any, any>, Base extends abstract new (...args: any[]) => any>(pipeFn: () => P, base: Base) {\n\tconst pipe = pipeFn()\n\tv.compile(pipe)\n\n\tabstract class Configurable extends (base as unknown as new (...args: any[]) => any) {\n\t\tdeclare static readonly Config: PipeOutput<P>\n\n\t\tprotected constructor (protected readonly config: PipeOutput<P>, ...baseArgs: BaseCtorParams<Base>) {\n\t\t\t// eslint-disable-next-line constructor-super\n\t\t\tsuper(...baseArgs)\n\t\t}\n\n\t\tstatic create<This extends Function & { prototype: any }>(\n\t\t\tthis: This,\n\t\t\tinput: ConditionalObjectKeys<PipeInput<P>>,\n\t\t\t...args: CtorParams<This> extends [PipeOutput<P>, ...infer R] ? R : never\n\t\t): This['prototype'] {\n\t\t\tconst r = v.validate(pipe, input)\n\t\t\tif (!r.valid) throw r.error\n\t\t\treturn new (this as any)(r.value, ...args) as This['prototype']\n\t\t}\n\t}\n\n\treturn Configurable as unknown as (abstract new (validated: PipeOutput<P>, ...baseArgs: BaseCtorParams<Base>) => InstanceType<Base> & { readonly config: PipeOutput<P> }) & {\n\t\treadonly Config: PipeOutput<P>\n\t\tcreate<This extends Function & { prototype: any }>(\n\t\t\tthis: This,\n\t\t\tinput: ConditionalObjectKeys<PipeInput<P>>,\n\t\t\t...args: CtorParams<This> extends [PipeOutput<P>, ...infer R] ? R : never\n\t\t): This['prototype']\n\t}\n}\n\nif (import.meta.vitest) {\n\tconst { describe, test, expect, expectTypeOf } = import.meta.vitest\n\tconst { v } = await import('valleyed')\n\n\tconst testPipe = () =>\n\t\tv.object({\n\t\t\thost: v.string(),\n\t\t\tport: v.number(),\n\t\t})\n\n\ttype TestConfig = PipeOutput<ReturnType<typeof testPipe>>\n\n\tclass TestBase {\n\t\tbaseValue: string\n\t\tconstructor() {\n\t\t\tthis.baseValue = 'base'\n\t\t}\n\t}\n\n\tclass TestBaseWithArgs {\n\t\tlabel: string\n\t\tconstructor(label: string) {\n\t\t\tthis.label = label\n\t\t}\n\t}\n\n\tdescribe('configurable', () => {\n\t\ttest('validation runs in static create before constructor body executes', () => {\n\t\t\tlet constructorRan = false\n\n\t\t\tconst Wrapped = configurable(testPipe, TestBase)\n\t\t\tclass MyClass extends Wrapped {\n\t\t\t\tprotected constructor(config: typeof MyClass.Config) {\n\t\t\t\t\tsuper(config)\n\t\t\t\t\tconstructorRan = true\n\t\t\t\t}\n\t\t\t}\n\n\t\t\texpect(() => MyClass.create({ host: 123, port: 'bad' } as any)).toThrow()\n\t\t\texpect(constructorRan).toBe(false)\n\n\t\t\tMyClass.create({ host: 'localhost', port: 3000 })\n\t\t\texpect(constructorRan).toBe(true)\n\t\t})\n\n\t\ttest('constructor receives validated value', () => {\n\t\t\tconst Wrapped = configurable(testPipe, TestBase)\n\t\t\tlet receivedConfig: unknown\n\n\t\t\tclass MyClass extends Wrapped {\n\t\t\t\tprotected constructor(config: typeof MyClass.Config) {\n\t\t\t\t\tsuper(config)\n\t\t\t\t\treceivedConfig = this.config\n\t\t\t\t}\n\t\t\t}\n\n\t\t\tMyClass.create({ host: 'localhost', port: 3000 })\n\n\t\t\texpect(receivedConfig).toEqual({ host: 'localhost', port: 3000 })\n\t\t})\n\n\t\ttest('external new is a compile error', () => {\n\t\t\tconst Wrapped = configurable(testPipe, TestBase)\n\t\t\tclass MyClass extends Wrapped {\n\t\t\t\tprotected constructor(config: typeof MyClass.Config) {\n\t\t\t\t\tsuper(config)\n\t\t\t\t}\n\t\t\t}\n\n\t\t\t// @ts-expect-error — external `new` on a class with protected constructor is a compile error\n\t\t\tvoid (() => new MyClass({ host: 'localhost', port: 3000 }))\n\t\t})\n\n\t\ttest('static Config resolves to PipeOutput<P> at the type level', () => {\n\t\t\tconst Wrapped = configurable(testPipe, TestBase)\n\t\t\tclass MyClass extends Wrapped {\n\t\t\t\tprotected constructor(config: typeof MyClass.Config) {\n\t\t\t\t\tsuper(config)\n\t\t\t\t}\n\t\t\t}\n\n\t\t\texpectTypeOf<typeof MyClass.Config>().toEqualTypeOf<TestConfig>()\n\t\t\texpect(MyClass.create).toBeTypeOf('function')\n\t\t})\n\n\t\ttest('ConstructorParameters<This>-based extras inference works for non-zero-arg leaf signatures', () => {\n\t\t\tconst Wrapped = configurable(testPipe, TestBase)\n\t\t\tclass MyClass extends Wrapped {\n\t\t\t\textra: number\n\t\t\t\tprotected constructor(config: typeof MyClass.Config, extra: number) {\n\t\t\t\t\tsuper(config)\n\t\t\t\t\tthis.extra = extra\n\t\t\t\t}\n\t\t\t}\n\n\t\t\tconst instance = MyClass.create({ host: 'localhost', port: 3000 }, 42)\n\n\t\t\texpect(instance.extra).toBe(42)\n\t\t\texpectTypeOf(instance).toHaveProperty('extra')\n\t\t\texpectTypeOf(instance.extra).toEqualTypeOf<number>()\n\n\t\t\t// @ts-expect-error — wrong extra type\n\t\t\tvoid (() => MyClass.create({ host: 'localhost', port: 3000 }, 'not-a-number'))\n\t\t})\n\n\t\ttest('base-args forwarding works for non-zero-arg bases', () => {\n\t\t\tconst Wrapped = configurable(testPipe, TestBaseWithArgs)\n\t\t\tclass MyClass extends Wrapped {\n\t\t\t\tprotected constructor(config: typeof MyClass.Config, label: string) {\n\t\t\t\t\tsuper(config, label)\n\t\t\t\t}\n\t\t\t}\n\n\t\t\tconst instance = MyClass.create({ host: 'localhost', port: 3000 }, 'test-label')\n\n\t\t\texpect(instance.label).toBe('test-label')\n\t\t})\n\n\t\ttest('config is accessible as protected readonly on instances', () => {\n\t\t\tconst Wrapped = configurable(testPipe, TestBase)\n\t\t\tclass MyClass extends Wrapped {\n\t\t\t\tprotected constructor(config: typeof MyClass.Config) {\n\t\t\t\t\tsuper(config)\n\t\t\t\t}\n\n\t\t\t\tgetHost() {\n\t\t\t\t\treturn this.config.host\n\t\t\t\t}\n\t\t\t}\n\n\t\t\tconst instance = MyClass.create({ host: 'localhost', port: 3000 })\n\t\t\texpect(instance.getHost()).toBe('localhost')\n\t\t})\n\n\t\ttest('accepts an abstract base class', () => {\n\t\t\tabstract class AbstractBase {\n\t\t\t\tabstract greet(): string\n\t\t\t}\n\n\t\t\tconst Wrapped = configurable(testPipe, AbstractBase)\n\t\t\tclass Concrete extends Wrapped {\n\t\t\t\tprotected constructor(config: typeof Concrete.Config) {\n\t\t\t\t\tsuper(config)\n\t\t\t\t}\n\t\t\t\tgreet() {\n\t\t\t\t\treturn `hello from ${this.config.host}`\n\t\t\t\t}\n\t\t\t}\n\n\t\t\tconst instance = Concrete.create({ host: 'localhost', port: 3000 })\n\t\t\texpect(instance.greet()).toBe('hello from localhost')\n\t\t\texpectTypeOf(instance).toHaveProperty('config')\n\t\t})\n\t})\n}\n","export const parseJSONValue = (data: any) => {\n\ttry {\n\t\tif (data?.constructor?.name !== 'String') return data\n\t\treturn JSON.parse(data)\n\t} catch {\n\t\treturn data\n\t}\n}\n","import crypto from 'crypto'\n\nexport function string(length = 20) {\n\treturn crypto.randomBytes(length).toString('hex').slice(0, length)\n}\n\nexport function number(min = 0, max = 2 ** 48 - 1) {\n\treturn crypto.randomInt(min, max)\n}\n","import { EquippedError } from '../errors'\n\nexport const sleep = (ms: number): Promise<void> => new Promise((resolve) => setTimeout(resolve, ms))\n\nexport const retry = async <T>(cb: () => Promise<{ done: true; value: T } | { done: false }>, tries: number, waitTimeInMs: number) => {\n\tif (tries <= 0) throw new EquippedError('out of tries', { tries, waitTimeInMs })\n\tconst result = await cb()\n\tif (result.done === true) return result.value\n\tawait sleep(waitTimeInMs)\n\treturn await retry(cb, tries - 1, waitTimeInMs)\n}\n","import axios from 'axios'\nimport { v } from 'valleyed'\n\nimport * as core from './core'\nimport { dbChangeConfigPipe, type DbChangeConfig } from './types'\nimport { EquippedError } from '../../../errors'\n\nexport const TopicPrefix = 'db-changes'\n\nexport abstract class DbChange<Model extends core.Model<core.IdType>, Entity extends core.Entity> {\n\tconstructor(\n\t\tprotected config: DbChangeConfig,\n\t\tprotected callbacks: core.DbChangeCallbacks<Model, Entity>,\n\t\tprotected mapper: (model: Model) => Entity,\n\t) {\n\t\tthis.config = v.assert(dbChangeConfigPipe(), config)\n\t}\n\n\tprotected async configureConnector(key: string, data: Record<string, string>) {\n\t\tconst instance = axios.create({ baseURL: this.config.debeziumUrl })\n\t\treturn await instance\n\t\t\t.put(`/connectors/${key}/config`, {\n\t\t\t\t'topic.prefix': TopicPrefix,\n\t\t\t\t'topic.creation.enable': 'false',\n\t\t\t\t'topic.creation.default.replication.factor': `-1`,\n\t\t\t\t'topic.creation.default.partitions': '-1',\n\t\t\t\t'key.converter': 'org.apache.kafka.connect.json.JsonConverter',\n\t\t\t\t'key.converter.schemas.enable': 'false',\n\t\t\t\t'value.converter': 'org.apache.kafka.connect.json.JsonConverter',\n\t\t\t\t'value.converter.schemas.enable': 'false',\n\t\t\t\t...data,\n\t\t\t})\n\t\t\t.then(async () => {\n\t\t\t\tconst topics = await instance.get(`/connectors/${key}/topics`)\n\t\t\t\treturn topics.data[key]?.topics?.includes?.(key) ?? false\n\t\t\t})\n\t\t\t.catch((err) => {\n\t\t\t\tthrow new EquippedError(`Failed to configure watcher`, { key }, err)\n\t\t\t})\n\t}\n}\n","import { v, type PipeOutput } from 'valleyed'\n\nimport { KafkaEventBus } from '../../../events/adapters/kafka'\n\nexport const dbChangeConfigPipe = () =>\n\tv.object({\n\t\tdebeziumUrl: v.string(),\n\t\teventBus: v.instanceOf(KafkaEventBus as unknown as abstract new (...args: any[]) => KafkaEventBus),\n\t})\n\nexport type DbChangeConfig = PipeOutput<ReturnType<typeof dbChangeConfigPipe>>\n\nexport type DbConfig = {\n\tchanges?: DbChangeConfig\n}\n","import { KafkaJS } from '@confluentinc/kafka-javascript'\nimport { v } from 'valleyed'\n\nimport { EquippedError } from '../../../errors'\nimport { Instance } from '../../../instance'\nimport type { Events } from '../../../types'\nimport { Random, configurable, parseJSONValue } from '../../../utilities'\nimport { EventBus, type Stream, type StreamOptions } from '../base'\n\nexport const kafkaConfigPipe = () =>\n\tv.meta(\n\t\tv.object({\n\t\t\tbrokers: v.array(v.string()),\n\t\t\tssl: v.optional(v.boolean()),\n\t\t\tsasl: v.optional(\n\t\t\t\tv.object({\n\t\t\t\t\tmechanism: v.is('plain' as const),\n\t\t\t\t\tusername: v.string(),\n\t\t\t\t\tpassword: v.string(),\n\t\t\t\t}),\n\t\t\t),\n\t\t\tclientId: v.optional(v.string()),\n\t\t}),\n\t\t{ title: 'Kafka Config', $refId: 'KafkaConfig' },\n\t)\n\nexport class KafkaEventBus extends configurable(kafkaConfigPipe, EventBus) {\n\t#client: KafkaJS.Kafka\n\t#admin: Promise<KafkaJS.Admin> | null = null\n\n\tprotected constructor(config: typeof KafkaEventBus.Config) {\n\t\tsuper(config)\n\t\tthis.#client = new KafkaJS.Kafka({\n\t\t\tkafkaJS: { ...config, logLevel: KafkaJS.logLevel.NOTHING },\n\t\t})\n\t}\n\n\tasync #getAdmin() {\n\t\tif (!this.#admin)\n\t\t\tthis.#admin = (async () => {\n\t\t\t\tconst admin = this.#client.admin()\n\t\t\t\tawait admin.connect()\n\t\t\t\treturn admin\n\t\t\t})()\n\t\treturn this.#admin\n\t}\n\n\tasync #createTopic(topic: string) {\n\t\tconst admin = await this.#getAdmin()\n\t\tawait admin.createTopics({ topics: [{ topic }], timeout: 5000 })\n\t}\n\n\tasync #deleteGroup(groupId: string) {\n\t\tconst admin = await this.#getAdmin()\n\t\tawait admin.deleteGroups([groupId]).catch(() => {})\n\t}\n\n\tstream<Event extends Events[keyof Events]>(topicName: Event['topic'], options: Partial<StreamOptions> = {}): Stream<Event['data']> {\n\t\tconst topic = options.skipScope ? topicName : Instance.get().getScopedName(topicName)\n\t\treturn {\n\t\t\tpublish: async (data) => {\n\t\t\t\tconst producer = this.#client.producer()\n\t\t\t\tawait producer.connect()\n\t\t\t\tawait producer.send({\n\t\t\t\t\ttopic,\n\t\t\t\t\tmessages: [{ value: JSON.stringify(data) }],\n\t\t\t\t})\n\t\t\t\treturn true\n\t\t\t},\n\t\t\tsubscribe: (onMessage) => {\n\t\t\t\tconst subscribe = async () => {\n\t\t\t\t\tawait this.#createTopic(topic)\n\t\t\t\t\tconst groupId = options.fanout\n\t\t\t\t\t\t? Instance.get().getScopedName(`${Instance.get().id}-fanout-${Random.string(10)}`)\n\t\t\t\t\t\t: topic\n\t\t\t\t\tconst consumer = this.#client.consumer({ kafkaJS: { groupId } })\n\n\t\t\t\t\tawait consumer.connect()\n\t\t\t\t\tawait consumer.subscribe({ topic })\n\n\t\t\t\t\tawait consumer.run({\n\t\t\t\t\t\teachMessage: async ({ message }) => {\n\t\t\t\t\t\t\tawait Instance.resolveBeforeCrash(async () => {\n\t\t\t\t\t\t\t\tif (!message.value) return\n\t\t\t\t\t\t\t\tawait onMessage(parseJSONValue(message.value.toString()))\n\t\t\t\t\t\t\t}).catch((error) =>\n\t\t\t\t\t\t\t\tInstance.crash(new EquippedError('Error processing kafka event', { topic, groupId, options }, error)),\n\t\t\t\t\t\t\t)\n\t\t\t\t\t\t},\n\t\t\t\t\t})\n\n\t\t\t\t\tif (options.fanout)\n\t\t\t\t\t\tInstance.on('close', async () => {\n\t\t\t\t\t\t\tawait consumer.disconnect()\n\t\t\t\t\t\t\tawait this.#deleteGroup(groupId)\n\t\t\t\t\t\t}, { class: KafkaEventBus })\n\t\t\t\t}\n\t\t\t\tInstance.on('start', subscribe, { class: KafkaEventBus })\n\t\t\t},\n\t\t}\n\t}\n}\n","import { AsyncLocalStorage } from 'node:async_hooks'\n\nimport {\n\tClientSession,\n\tCollection,\n\tMongoClient,\n\tObjectId,\n\ttype CollectionInfo,\n\ttype OptionalUnlessRequiredId,\n\ttype SortDirection,\n\ttype WithId,\n} from 'mongodb'\nimport { v } from 'valleyed'\n\nimport { MongoDbChange } from './changes'\nimport { mongoDbConfigPipe, type MongoDbConfig } from './pipes'\nimport { parseMongodbQueryParams } from './query'\nimport { EquippedError } from '../../../errors'\nimport { Instance } from '../../../instance'\nimport type { QueryParamsBase } from '../../pipes'\nimport * as core from '../base/core'\nimport { Db } from '../base/db'\nimport type { DbConfig } from '../base/types'\n\nconst idKey = '_id'\ntype IdType = { _id: string }\n\nconst sessionStore = new AsyncLocalStorage<ClientSession | undefined>(undefined)\n\nexport class MongoDb extends Db<{ _id: string }> {\n\tclient: MongoClient\n\t#cols: { db: string; col: string }[] = []\n\n\tconstructor(\n\t\tprivate mongoConfig: MongoDbConfig,\n\t\tdbConfig: DbConfig,\n\t) {\n\t\tsuper(dbConfig)\n\t\tthis.mongoConfig = v.assert(mongoDbConfigPipe(), mongoConfig)\n\t\tthis.client = new MongoClient(mongoConfig.uri, { ignoreUndefined: true })\n\t\tInstance.on('start', async () => {\n\t\t\tawait this.client.connect()\n\n\t\t\tconst grouped = this.#cols.reduce<Record<string, string[]>>((acc, cur) => {\n\t\t\t\tif (!acc[cur.db]) acc[cur.db] = []\n\t\t\t\tacc[cur.db].push(cur.col)\n\t\t\t\treturn acc\n\t\t\t}, {})\n\n\t\t\tconst options = {\n\t\t\t\tchangeStreamPreAndPostImages: { enabled: true },\n\t\t\t}\n\t\t\tawait Promise.all(\n\t\t\t\tObject.entries(grouped).map(async ([dbName, colNames]) => {\n\t\t\t\t\tconst db = this.client.db(dbName)\n\t\t\t\t\tconst collections = await db.listCollections<CollectionInfo>().toArray()\n\t\t\t\t\treturn colNames.map(async (colName) => {\n\t\t\t\t\t\tconst existing = collections.find((collection) => collection.name === colName)\n\t\t\t\t\t\tif (existing) {\n\t\t\t\t\t\t\tif (\n\t\t\t\t\t\t\t\texisting.options?.changeStreamPreAndPostImages?.enabled !== options.changeStreamPreAndPostImages.enabled\n\t\t\t\t\t\t\t)\n\t\t\t\t\t\t\t\tawait db.command({ collMod: colName, ...options })\n\t\t\t\t\t\t} else await db.createCollection(colName, options)\n\t\t\t\t\t})\n\t\t\t\t}),\n\t\t\t)\n\t\t})\n\t\tInstance.on('close', async () => this.client.close())\n\t}\n\n\tasync session<T>(callback: () => Promise<T>) {\n\t\tif (sessionStore.getStore()) return callback()\n\t\tconst session = await this.client.startSession()\n\t\treturn session.withTransaction(async () => sessionStore.run(session, callback))\n\t}\n\n\tid() {\n\t\treturn new ObjectId()\n\t}\n\n\tuse<Model extends core.Model<{ _id: string }>, Entity extends core.Entity>(config: core.Config<Model, Entity>) {\n\t\tconst db = this.getScopedDb(config.db)\n\t\tthis.#cols.push({ db, col: config.col })\n\t\treturn this.#getTable(config, this.client.db(db).collection<Model>(config.col))\n\t}\n\n\t#getTable<Model extends core.Model<IdType>, Entity extends core.Entity>(\n\t\tconfig: core.Config<Model, Entity>,\n\t\tcollection: Collection<Model>,\n\t) {\n\t\ttype WI = Model | WithId<Model>\n\t\tfunction transform(doc: WI, select?: core.Select<Entity>): Entity\n\t\t// eslint-disable-next-line no-redeclare\n\t\tfunction transform(doc: WI[], select?: core.Select<Entity>): Entity[]\n\t\t// eslint-disable-next-line no-redeclare\n\t\tfunction transform(doc: WI | WI[], select?: core.Select<Entity>) {\n\t\t\tconst docs = Array.isArray(doc) ? doc : [doc]\n\t\t\tconst mapped = docs.map((d) => config.mapper(d as Model, select))\n\t\t\treturn Array.isArray(doc) ? mapped : mapped[0]\n\t\t}\n\n\t\tfunction prepInsertValue(value: core.CreateInput<Model>, id: string, now: Date, skipUpdate?: boolean) {\n\t\t\tconst base: core.Model<IdType> = {\n\t\t\t\t[idKey]: id,\n\t\t\t\t...(config.options?.skipAudit\n\t\t\t\t\t? {}\n\t\t\t\t\t: {\n\t\t\t\t\t\t\tcreatedAt: now.getTime(),\n\t\t\t\t\t\t\t...(skipUpdate ? {} : { updatedAt: now.getTime() }),\n\t\t\t\t\t\t}),\n\t\t\t}\n\t\t\treturn {\n\t\t\t\t...value,\n\t\t\t\t...base,\n\t\t\t} as unknown as OptionalUnlessRequiredId<Model>\n\t\t}\n\n\t\tfunction prepUpdateValue(value: core.UpdateInput<Model>, now: Date, upsert = false) {\n\t\t\treturn {\n\t\t\t\t...value,\n\t\t\t\t$set: {\n\t\t\t\t\t...value.$set,\n\t\t\t\t\t...(upsert || (Object.keys(value).length > 0 && !config.options?.skipAudit) ? { updatedAt: now.getTime() } : {}),\n\t\t\t\t},\n\t\t\t}\n\t\t}\n\n\t\tconst dbThis = this\n\n\t\tconst table: core.Table<IdType, Model, Entity, { collection: Collection<Model> }> = {\n\t\t\tconfig,\n\t\t\textras: { collection },\n\n\t\t\tquery: async (params) => {\n\t\t\t\tconst results = await parseMongodbQueryParams(collection, params as QueryParamsBase)\n\t\t\t\treturn {\n\t\t\t\t\t...results,\n\t\t\t\t\tresults: transform(results.results as any, params.select as core.Select<Entity>),\n\t\t\t\t} as any\n\t\t\t},\n\n\t\t\tfindMany: async (filter, options = {}) => {\n\t\t\t\tconst { select, ...rest } = options\n\t\t\t\tconst projection = select?.length ? Object.fromEntries(select.map((k) => [k, 1])) : undefined\n\t\t\t\tconst sortArray = Array.isArray(rest.sort) ? rest.sort : rest.sort ? [rest.sort] : []\n\t\t\t\tconst sort = sortArray.map((p) => [p.field, p.desc ? 'desc' : 'asc'] as [string, SortDirection])\n\t\t\t\tconst docs = await collection\n\t\t\t\t\t.find(filter, {\n\t\t\t\t\t\tsession: sessionStore.getStore(),\n\t\t\t\t\t\tlimit: rest.limit,\n\t\t\t\t\t\tsort,\n\t\t\t\t\t\tprojection,\n\t\t\t\t\t})\n\t\t\t\t\t.toArray()\n\t\t\t\treturn transform(docs, select as core.Select<Entity>) as any\n\t\t\t},\n\n\t\t\tfindOne: async (filter: core.Filter<Model>, options = {}) => {\n\t\t\t\tconst result = await table.findMany(filter, { ...options, limit: 1 } as any)\n\t\t\t\treturn (result.at(0) ?? null) as any\n\t\t\t},\n\n\t\t\tfindById: async (id, options = {}) => {\n\t\t\t\tconst result = await table.findOne({ [idKey]: id } as core.Filter<Model>, options as any)\n\t\t\t\treturn result as any\n\t\t\t},\n\n\t\t\tinsertMany: async (values, options = {}) => {\n\t\t\t\tconst now = options.getTime?.() ?? new Date()\n\t\t\t\tconst payload = values.map((value, i) => prepInsertValue(value, options.makeId?.(i) ?? new ObjectId().toString(), now))\n\t\t\t\tawait collection.insertMany(payload, { session: sessionStore.getStore() })\n\n\t\t\t\tconst insertedData = await Promise.all(payload.map(async (data) => await table.findById(data[idKey] as any)))\n\t\t\t\treturn insertedData.filter((value) => !!value)\n\t\t\t},\n\n\t\t\tinsertOne: async (values, options = {}) => {\n\t\t\t\tconst result = await table.insertMany([values], options)\n\t\t\t\treturn result[0]\n\t\t\t},\n\n\t\t\tupdateMany: async (filter, values, options = {}) => {\n\t\t\t\tconst now = options.getTime?.() ?? new Date()\n\t\t\t\tconst session = sessionStore.getStore()\n\t\t\t\tconst data = await collection.find(filter, { session, projection: { [idKey]: 1 } }).toArray()\n\t\t\t\tconst ids = data.map((doc) => doc[idKey])\n\t\t\t\tconst filterUpd = { [idKey]: { $in: ids } } as core.Filter<Model>\n\t\t\t\tawait collection.updateMany(filterUpd, prepUpdateValue(values, now), { session })\n\t\t\t\treturn table.findMany(filterUpd)\n\t\t\t},\n\n\t\t\tupdateOne: async (filter, values, options = {}) => {\n\t\t\t\tconst now = options.getTime?.() ?? new Date()\n\t\t\t\tconst doc = await collection.findOneAndUpdate(filter, prepUpdateValue(values, now), {\n\t\t\t\t\treturnDocument: 'after',\n\t\t\t\t\tsession: sessionStore.getStore(),\n\t\t\t\t})\n\t\t\t\treturn doc ? transform(doc, undefined) : null\n\t\t\t},\n\n\t\t\tupdateById: async (id, values, options = {}) => {\n\t\t\t\tconst result = await table.updateOne({ [idKey]: id } as core.Filter<Model>, values, options)\n\t\t\t\treturn result\n\t\t\t},\n\n\t\t\tupsertOne: async (filter, values, options = {}) => {\n\t\t\t\tconst now = options.getTime?.() ?? new Date()\n\n\t\t\t\tconst doc = await collection.findOneAndUpdate(\n\t\t\t\t\tfilter,\n\t\t\t\t\t{\n\t\t\t\t\t\t...prepUpdateValue('update' in values ? values.update : {}, now, true),\n\t\t\t\t\t\t// @ts-expect-error fighting ts\n\t\t\t\t\t\t$setOnInsert: prepInsertValue(values.insert, options.makeId?.() ?? new ObjectId().toString(), now, true),\n\t\t\t\t\t},\n\t\t\t\t\t{ returnDocument: 'after', session: sessionStore.getStore(), upsert: true },\n\t\t\t\t)\n\n\t\t\t\treturn transform(doc, undefined)\n\t\t\t},\n\n\t\t\tdeleteMany: async (filter, options = {}) => {\n\t\t\t\tconst docs = await table.findMany(filter, options)\n\t\t\t\tawait collection.deleteMany(filter, { session: sessionStore.getStore() })\n\t\t\t\treturn docs\n\t\t\t},\n\n\t\t\tdeleteOne: async (filter) => {\n\t\t\t\tconst doc = await collection.findOneAndDelete(filter, { session: sessionStore.getStore() })\n\t\t\t\treturn doc ? transform(doc, undefined) : null\n\t\t\t},\n\n\t\t\tdeleteById: async (id) => {\n\t\t\t\tconst result = await table.deleteOne({ [idKey]: id } as core.Filter<Model>)\n\t\t\t\treturn result\n\t\t\t},\n\n\t\t\tbulkWrite: async (operations, options = {}) => {\n\t\t\t\tif (!operations.length) return\n\t\t\t\tconst bulk = collection.initializeUnorderedBulkOp({ session: sessionStore.getStore() })\n\t\t\t\tconst now = options.getTime?.() ?? new Date()\n\t\t\t\toperations.forEach((operation, i) => {\n\t\t\t\t\tswitch (operation.op) {\n\t\t\t\t\t\tcase 'insert':\n\t\t\t\t\t\t\tbulk.insert(prepInsertValue(operation.value, operation.makeId?.(i) ?? new ObjectId().toString(), now))\n\t\t\t\t\t\t\tbreak\n\t\t\t\t\t\tcase 'delete':\n\t\t\t\t\t\t\tbulk.find(operation.filter).delete()\n\t\t\t\t\t\t\tbreak\n\t\t\t\t\t\tcase 'update':\n\t\t\t\t\t\t\tbulk.find(operation.filter).update(prepUpdateValue(operation.value, now))\n\t\t\t\t\t\t\tbreak\n\t\t\t\t\t\tcase 'upsert':\n\t\t\t\t\t\t\tbulk.find(operation.filter)\n\t\t\t\t\t\t\t\t.upsert()\n\t\t\t\t\t\t\t\t.update({\n\t\t\t\t\t\t\t\t\t...prepUpdateValue('update' in operation ? operation.update : {}, now, true),\n\t\t\t\t\t\t\t\t\t$setOnInsert: prepInsertValue(\n\t\t\t\t\t\t\t\t\t\toperation.insert as any,\n\t\t\t\t\t\t\t\t\t\toperation.makeId?.(i) ?? new ObjectId().toString(),\n\t\t\t\t\t\t\t\t\t\tnow,\n\t\t\t\t\t\t\t\t\t\ttrue,\n\t\t\t\t\t\t\t\t\t),\n\t\t\t\t\t\t\t\t})\n\t\t\t\t\t\t\tbreak\n\t\t\t\t\t\tdefault:\n\t\t\t\t\t\t\tthrow new EquippedError(`Unknown bulkWrite operation`, { operation })\n\t\t\t\t\t}\n\t\t\t\t})\n\t\t\t\tawait bulk.execute({ session: sessionStore.getStore() })\n\t\t\t},\n\n\t\t\twatch(callbacks) {\n\t\t\t\tif (!dbThis.config.changes)\n\t\t\t\t\tInstance.crash(new EquippedError('Db changes are not enabled in the configuration.', { config }))\n\t\t\t\treturn new MongoDbChange<Model, Entity>(dbThis.mongoConfig, dbThis.config.changes, collection, callbacks, (m) =>\n\t\t\t\t\tconfig.mapper(m),\n\t\t\t\t)\n\t\t\t},\n\t\t}\n\n\t\treturn table\n\t}\n}\n","import { type PipeInput, v } from 'valleyed'\n\nexport const mongoDbConfigPipe = () =>\n\tv.meta(\n\t\tv.object({\n\t\t\turi: v.string(),\n\t\t}),\n\t\t{ title: 'Mongodb Config', $refId: 'MongodbConfig' },\n\t)\n\nexport type MongoDbConfig = PipeInput<ReturnType<typeof mongoDbConfigPipe>>\n","import { Collection } from 'mongodb'\n\nimport { QueryKeys, type QueryParamsBase, type QueryResults, type QueryWhereBlock, type QueryWhereClause } from '../../pipes'\nimport * as core from '../base/core'\n\nexport const parseMongodbQueryParams = async <Model extends core.Model<{ _id: string }>>(\n\tcollection: Collection<Model>,\n\tparams: QueryParamsBase,\n): Promise<QueryResults<Model>> => {\n\t// Handle where/search clauses\n\tconst query = <ReturnType<typeof buildWhereQuery>[]>[]\n\tconst where = buildWhereQuery(params.where, params.whereType)\n\tif (where) query.push(where)\n\tconst auth = buildWhereQuery(params.auth, params.authType)\n\tif (auth) query.push(auth)\n\tif (params.search && params.search.fields.length > 0) {\n\t\tconst search = params.search.fields.map((field) => ({\n\t\t\t[field]: {\n\t\t\t\t$regex: new RegExp(params.search!.value, 'i'),\n\t\t\t},\n\t\t}))\n\t\tquery.push({ $or: search })\n\t}\n\tconst totalClause = {}\n\tif (query.length > 0) totalClause['$and'] = query\n\n\t// Handle sort clauses\n\tconst sort = params.sort.map((p) => [p.field, p.desc ? 'desc' : 'asc'])\n\n\t// Handle limit/offest clause\n\tconst all = params.all ?? false\n\tconst limit = params.limit\n\tconst page = params.page\n\n\tconst total = await collection.countDocuments(totalClause)\n\n\tconst projection = params.select?.length ? Object.fromEntries(params.select.map((k) => [k, 1])) : undefined\n\tlet builtQuery = collection.find(totalClause, { projection })\n\tif (sort.length) builtQuery = builtQuery.sort(Object.fromEntries(sort))\n\tif (!all && limit) {\n\t\tbuiltQuery = builtQuery.limit(limit)\n\t\tif (page) builtQuery = builtQuery.skip((page - 1) * limit)\n\t}\n\n\tconst results = await builtQuery.toArray()\n\tconst start = 1\n\tconst last = Math.ceil(total / limit) || 1\n\tconst next = page >= last ? null : page + 1\n\tconst previous = page <= start ? null : page - 1\n\n\treturn {\n\t\tpages: { start, last, next, previous, current: page },\n\t\tdocs: { limit, total, count: results.length },\n\t\tresults: results as any[],\n\t} satisfies QueryResults<Model>\n}\n\nfunction isWhereBlock(param: QueryWhereClause): param is QueryWhereBlock {\n\treturn Object.values(QueryKeys).includes(param.condition as QueryKeys)\n}\n\nconst buildWhereQuery = (params: QueryWhereClause[], key: QueryKeys = QueryKeys.and): Record<string, Record<string, any>> | null => {\n\tconst where = (Array.isArray(params) ? params : [])\n\t\t.map((param) => {\n\t\t\tif (isWhereBlock(param)) return buildWhereQuery(param.value, param.condition)\n\t\t\tconst { field } = param\n\t\t\tconst checkedField = field === 'id' ? '_id' : (field ?? '')\n\t\t\tconst checkedValue = param.value === undefined ? '' : param.value\n\t\t\treturn {\n\t\t\t\tfield: checkedField,\n\t\t\t\tvalue: checkedValue,\n\t\t\t\tcondition: param.condition,\n\t\t\t\tisWhere: true,\n\t\t\t}\n\t\t})\n\t\t.filter((c) => !!c)\n\t\t.map((c) => {\n\t\t\tif (c.isWhere) return { [`${c.field}`]: { [`$${c.condition}`]: c.value } }\n\t\t\telse return c\n\t\t})\n\n\treturn where.length > 0 ? { [`$${key}`]: where } : null\n}\n","import { type ConditionalObjectKeys, type Pipe, type PipeInput, type PipeOutput, v } from 'valleyed'\n\nimport { Instance } from '../instance'\nimport type { Select } from './adapters/base/core'\n\nexport enum QueryKeys {\n\tand = 'and',\n\tor = 'or',\n}\n\nexport enum Conditions {\n\tlt = 'lt',\n\tlte = 'lte',\n\tgt = 'gt',\n\tgte = 'gte',\n\teq = 'eq',\n\tne = 'ne',\n\tin = 'in',\n\tnin = 'nin',\n}\n\n// eslint-disable-next-line promise/valid-params -- valleyed v.catch(schema, fallback) is not Promise.prototype.catch\nconst queryKeys = v.catch(v.defaults(v.in([QueryKeys.and, QueryKeys.or]), QueryKeys.and), QueryKeys.and)\nconst queryWhere = v.object({\n\tfield: v.string(),\n\tvalue: v.any(),\n\t// eslint-disable-next-line promise/valid-params -- valleyed v.catch(schema, fallback) is not Promise.prototype.catch\n\tcondition: v.catch(v.defaults(v.in(Object.values(Conditions)), Conditions.eq), Conditions.eq),\n})\nconst queryWhereBlock = v.recursive(\n\t() =>\n\t\tv.discriminate((d) => (Object.values(QueryKeys).includes(d.condition as any) ? 'block' : 'regular'), {\n\t\t\tblock: v.object({\n\t\t\t\tcondition: queryKeys,\n\t\t\t\tvalue: v.array(queryWhereBlock),\n\t\t\t}),\n\t\t\tregular: queryWhere,\n\t\t}),\n\t'QueryWhereBlock',\n) as Pipe<\n\t| {\n\t\t\tcondition: PipeInput<typeof queryKeys>\n\t\t\tvalue: PipeInput<typeof queryWhere>[]\n\t  }\n\t| PipeInput<typeof queryWhere>,\n\t| {\n\t\t\tcondition: PipeOutput<typeof queryKeys>\n\t\t\tvalue: PipeOutput<typeof queryWhere>[]\n\t  }\n\t| PipeOutput<typeof queryWhere>\n>\n\nconst queryWhereClause = v.defaults(v.array(queryWhereBlock), [])\n\nexport function queryParamsPipe() {\n\treturn v.meta(\n\t\tv\n\t\t\t.object({\n\t\t\t\tall: v.defaults(v.boolean(), false),\n\t\t\t\tlimit: v.lazy(() => {\n\t\t\t\t\tconst pagLimit = Instance.get().settings.utils.paginationDefaultLimit\n\t\t\t\t\t// eslint-disable-next-line promise/valid-params -- valleyed v.catch(schema, fallback) is not Promise.prototype.catch\n\t\t\t\t\treturn v.catch(v.defaults(v.number().pipe(v.lte(pagLimit)), pagLimit), pagLimit)\n\t\t\t\t}),\n\t\t\t\t// eslint-disable-next-line promise/valid-params -- valleyed v.catch(schema, fallback) is not Promise.prototype.catch\n\t\t\t\tpage: v.catch(v.defaults(v.number().pipe(v.gte(1)), 1), 1),\n\t\t\t\tsearch: v.defaults(\n\t\t\t\t\tv.nullish(\n\t\t\t\t\t\tv.object({\n\t\t\t\t\t\t\tvalue: v.string(),\n\t\t\t\t\t\t\tfields: v.array(v.string()),\n\t\t\t\t\t\t}),\n\t\t\t\t\t),\n\t\t\t\t\tnull,\n\t\t\t\t),\n\t\t\t\tsort: v.defaults(\n\t\t\t\t\tv.array(\n\t\t\t\t\t\tv.object({\n\t\t\t\t\t\t\tfield: v.string(),\n\t\t\t\t\t\t\tdesc: v.defaults(v.boolean(), false),\n\t\t\t\t\t\t}),\n\t\t\t\t\t),\n\t\t\t\t\t[],\n\t\t\t\t),\n\t\t\t\twhereType: queryKeys,\n\t\t\t\twhere: queryWhereClause,\n\t\t\t\tselect: v.optional(v.array(v.string())),\n\t\t\t})\n\t\t\t.pipe((p) => ({ ...p, auth: <(typeof p)['where']>[], authType: QueryKeys.and })),\n\t\t{ title: 'Query Params', $refId: 'QueryParams' },\n\t)\n}\n\nexport function queryResultsPipe<T>(model: Pipe<any, T>) {\n\treturn v.object({\n\t\tpages: v.object({\n\t\t\tcurrent: v.number(),\n\t\t\tstart: v.number(),\n\t\t\tlast: v.number(),\n\t\t\tprevious: v.nullable(v.number()),\n\t\t\tnext: v.nullable(v.number()),\n\t\t}),\n\t\tdocs: v.object({\n\t\t\tlimit: v.number(),\n\t\t\ttotal: v.number(),\n\t\t\tcount: v.number(),\n\t\t}),\n\t\tresults: v.array(model),\n\t})\n}\n\nexport function wrapQueryParams(params: QueryParamsInput): QueryParamsBase {\n\treturn v.assert(queryParamsPipe(), params)\n}\n\nexport type QueryParamsBase = PipeOutput<ReturnType<typeof queryParamsPipe>>\nexport type QueryParams<T, S extends Select<T> | undefined = undefined> = Omit<QueryParamsBase, 'select'> & {\n\tselect?: S extends undefined ? QueryParamsBase['select'] : S\n}\nexport type QueryParamsInput = ConditionalObjectKeys<PipeInput<ReturnType<typeof queryParamsPipe>>>\nexport type QueryWhereClause = QueryParamsBase['where'][number]\nexport type QueryWhere = Extract<QueryWhereClause, { field: string }>\nexport type QueryWhereBlock = Exclude<QueryWhereClause, { field: string }>\nexport type QueryResults<T> = PipeOutput<ReturnType<typeof queryResultsPipe<T>>>\n","import * as core from './core'\nimport type { DbConfig } from './types'\nimport { Instance } from '../../../instance'\n\nexport type TableOptions = { skipAudit?: boolean }\n\nexport abstract class Db<IdKey extends core.IdType> {\n\tconstructor(protected config: DbConfig) {}\n\n\tprotected getScopedDb(db: string) {\n\t\treturn Instance.get().getScopedName(db).replaceAll('.', '-')\n\t}\n\n\tabstract use<Model extends core.Model<IdKey>, Entity extends core.Entity>(\n\t\tconfig: core.Config<Model, Entity>,\n\t): core.Table<IdKey, Model, Entity>\n\n\tabstract session<T>(callback: () => Promise<T>): Promise<T>\n}\n"],"mappings":"qkBAAA,IAAAA,GAAA,GAAAC,EAAAD,GAAA,aAAAE,EAAA,kBAAAC,EAAA,sBAAAC,IAAA,eAAAC,GAAAL,ICAA,IAAAM,GAAwC,mBACxCC,EAAuB,oBCDhB,IAAMC,EAAN,cAA4B,KAAM,CACxC,YACiBC,EACAC,EACAC,EACf,CACD,MAAMF,EAAS,CAAE,MAAAE,CAAM,CAAC,EAJR,aAAAF,EACA,aAAAC,EACA,WAAAC,CAGjB,CACD,ECCO,IAAeC,EAAf,KAAwB,CAE/B,ECXA,IAAAC,EAAkC,qBAClCC,GAAqB,gBACrBC,EAA6B,oBCQtB,SAASC,EAAaC,EAAqBC,EAA0B,CAC3E,QAAWC,KAAOD,EAAO,MACxB,GAAI,CAACD,EAAM,KAAMG,GAAMA,EAAE,QAAUD,CAAG,EAAG,CACxC,IAAME,EAAUF,EAAI,MAAQ,UACtBG,EAAYJ,EAAO,OAAO,MAAQ,YACxC,MAAM,IAAI,MAAM,uBAAuBI,CAAS,qBAAqBD,CAAO,UAAUA,CAAO,oBAAoB,CAClH,CAGD,GAAIH,EAAO,MAAO,CACjB,IAAMK,EAAY,CAAC,GAAGN,EAAOC,CAAM,EACnCM,GAAYD,CAAS,CACtB,CAEAN,EAAM,KAAKC,CAAM,CAClB,CAEA,SAASM,GAAYP,EAA2B,CAC/C,IAAMQ,EAAe,IAAI,IACzB,QAAWL,KAAKH,EAAO,CACtB,GAAI,CAACG,EAAE,MAAO,SACd,IAAMM,EAAWD,EAAa,IAAIL,EAAE,KAAK,EACzC,GAAIM,EACH,QAAWP,KAAOC,EAAE,MACdM,EAAS,SAASP,CAAG,GAAGO,EAAS,KAAKP,CAAG,OAG/CM,EAAa,IAAIL,EAAE,MAAO,CAAC,GAAGA,EAAE,KAAK,CAAC,CAExC,CAEA,IAAMO,EAAU,IAAI,IACdC,EAAQ,IAAI,IAElB,SAASC,EAAMC,EAAeC,EAAwB,CACrD,GAAIH,EAAM,IAAIE,CAAG,EAAG,CACnB,IAAME,EAAaD,EAAK,QAAQD,CAAG,EAE7BG,EADQ,CAAC,GAAGF,EAAK,MAAMC,CAAU,EAAGF,CAAG,EACzB,IAAKI,GAAMA,EAAE,MAAQ,SAAS,EAClD,MAAM,IAAI,MAAM,mBAAmBD,EAAM,KAAK,UAAK,CAAC,EAAE,CACvD,CACA,GAAI,CAAAN,EAAQ,IAAIG,CAAG,EACnB,CAAAF,EAAM,IAAIE,CAAG,EACbC,EAAK,KAAKD,CAAG,EACb,QAAWX,KAAOM,EAAa,IAAIK,CAAG,GAAK,CAAC,EAC3CD,EAAMV,EAAKY,CAAI,EAEhBA,EAAK,IAAI,EACTH,EAAM,OAAOE,CAAG,EAChBH,EAAQ,IAAIG,CAAG,EAChB,CAEA,QAAWA,KAAOL,EAAa,KAAK,EACnCI,EAAMC,EAAK,CAAC,CAAC,CAEf,CAIO,SAASK,GAAelB,EAAqBmB,EAAkB,GAAoB,CACzF,GAAInB,EAAM,SAAW,EAAG,MAAO,CAAC,EAEhC,IAAMoB,EAAe,IAAI,IACnBC,EAA0B,CAAC,EAEjC,QAAWlB,KAAKH,EACf,GAAIG,EAAE,MAAO,CACZ,IAAMmB,EAAOF,EAAa,IAAIjB,EAAE,KAAK,EACjCmB,EAAMA,EAAK,KAAKnB,CAAC,EAChBiB,EAAa,IAAIjB,EAAE,MAAO,CAACA,CAAC,CAAC,CACnC,MACCkB,EAAU,KAAKlB,CAAC,EAIlB,IAAMoB,EAAU,CAAC,GAAGH,EAAa,KAAK,CAAC,EAEjCI,EAAW,IAAI,IACfC,EAAQ,IAAI,IAElB,QAAWZ,KAAOU,EACjBC,EAAS,IAAIX,EAAK,CAAC,EACnBY,EAAM,IAAIZ,EAAK,CAAC,CAAC,EAGlB,QAAWV,KAAKH,EACf,GAAKG,EAAE,MACP,QAAWD,KAAOC,EAAE,MAAO,CAC1B,GAAM,CAACuB,EAAMC,CAAE,EAAIR,EAAS,CAAChB,EAAE,MAAOD,CAAG,EAAI,CAACA,EAAKC,EAAE,KAAK,EACpDyB,EAAQH,EAAM,IAAIC,CAAI,EACvBE,EAAM,SAASD,CAAE,IACrBC,EAAM,KAAKD,CAAE,EACbH,EAAS,IAAIG,EAAIH,EAAS,IAAIG,CAAE,EAAK,CAAC,EAExC,CAGD,IAAME,EAAsB,CAAC,EACzBC,EAAQP,EAAQ,OAAQN,GAAMO,EAAS,IAAIP,CAAC,IAAM,CAAC,EAEvD,KAAOa,EAAM,OAAS,GAAG,CACxB,IAAMC,EAAsB,CAAC,EACvBC,EAAmB,CAAC,EAE1B,QAAWnB,KAAOiB,EACjBC,EAAM,KAAK,GAAIX,EAAa,IAAIP,CAAG,GAAK,CAAC,CAAE,EAE5CgB,EAAO,KAAKE,CAAK,EAEjB,QAAWlB,KAAOiB,EACjB,QAAWG,KAAYR,EAAM,IAAIZ,CAAG,GAAK,CAAC,EAAG,CAC5C,IAAMqB,EAAIV,EAAS,IAAIS,CAAQ,EAAK,EACpCT,EAAS,IAAIS,EAAUC,CAAC,EACpBA,IAAM,GAAGF,EAAK,KAAKC,CAAQ,CAChC,CAGDH,EAAQE,CACT,CAEA,GAAIX,EAAU,OAAS,EAAG,CACzB,IAAMc,EAAed,EAAU,OAAQlB,GAAMA,EAAE,MAAM,OAAS,CAAC,EACzDiC,EAAkBf,EAAU,OAAQlB,GAAMA,EAAE,MAAM,SAAW,CAAC,EAEpE,GAAIgC,EAAa,OAAS,EAAG,CAC5B,IAAIE,EAAW,GACf,QAAWlC,KAAKgC,EACf,QAAWjC,KAAOC,EAAE,MAAO,CAC1B,IAAMmC,EAAWT,EAAO,UAAWU,GAAMA,EAAE,KAAMC,GAAMA,EAAE,QAAUtC,CAAG,CAAC,EACnEoC,EAAWD,IAAUA,EAAWC,EACrC,CAED,IAAMG,EAAWJ,EAAW,EACxBI,EAAWZ,EAAO,OACrBA,EAAOY,CAAQ,EAAE,KAAK,GAAGN,CAAY,EAErCN,EAAO,KAAKM,CAAY,CAE1B,CAEIC,EAAgB,OAAS,IACxBP,EAAO,OAAS,EACnBA,EAAOA,EAAO,OAAS,CAAC,EAAE,KAAK,GAAGO,CAAe,EAEjDP,EAAO,KAAKO,CAAe,EAG9B,CAEA,OAAOP,CACR,CAEA,eAAsBa,EACrB1C,EACA2C,EAAmCC,GAAU,CAC5C,MAAMA,CACP,EACAzB,EAAkB,GACjB,CACD,IAAMU,EAASX,GAAelB,EAAOmB,CAAM,EAC3C,QAAWY,KAASF,EACnB,MAAM,QAAQ,IACbE,EAAM,IAAI,MAAO5B,GAAM,CACtB,GAAI,CACH,OAAI,OAAOA,EAAE,IAAO,WAAmB,MAAMA,EAAE,GAAG,EAC3C,MAAMA,EAAE,EAChB,OAASyC,EAAO,CACf,OAAOD,EAAQC,aAAiB,MAAQA,EAAQ,IAAI,MAAM,GAAGA,CAAK,EAAE,CAAC,CACtE,CACD,CAAC,CACF,CACF,CCrLA,IAAAC,EAA+E,oBAElEC,EAAuB,IACnC,IAAE,OAAO,CACR,IAAK,IAAE,OAAO,CACb,KAAM,IAAE,OAAO,CAChB,CAAC,EACD,IAAK,IAAE,SACN,IAAE,OAAO,CACR,MAAO,IAAE,SAAS,IAAE,GAAG,CAAC,QAAS,QAAS,OAAQ,OAAQ,QAAS,QAAS,QAAQ,CAAU,EAAG,MAAM,CACxG,CAAC,EACD,CAAC,CACF,EACA,MAAO,IAAE,SACR,IAAE,OAAO,CACR,eAAgB,IAAE,SAAS,IAAE,OAAO,EAAG,EAAE,EACzC,uBAAwB,IAAE,SAAS,IAAE,OAAO,EAAG,GAAG,EAClD,sBAAuB,IAAE,SAAS,IAAE,OAAO,EAAG,GAAG,CAClD,CAAC,EACD,CAAC,CACF,CACD,CAAC,EFXK,IAAMC,EAAN,MAAMC,CAAS,CACrB,MAAOC,GACP,MAAOC,GACP,MAAOC,GAAmD,CAAC,EAClD,SACA,IAED,YAAYC,EAAoB,CACvCJ,EAASE,GAAY,KACrB,KAAK,SAAW,OAAO,OAAOE,CAAQ,EACtC,KAAK,OAAM,EAAAC,SAAY,CACtB,MAAO,KAAK,SAAS,IAAI,MACzB,YAAa,CACZ,IAAK,EAAAA,QAAK,eAAe,IACzB,MAAO,EAAAA,QAAK,eAAe,IAC3B,IAAK,EAAAA,QAAK,eAAe,IACzB,IAAK,EAAAA,QAAK,eAAe,GAC1B,EACA,MAAO,KAAO,CACb,WAAYL,EAASC,EACtB,EACD,CAAC,EACDD,EAASM,GAAuB,CACjC,CAEA,MAAMC,EAAY,CACjB,GAAIP,EAASC,KAAQ,OAAW,OAAOD,EAAS,MAAM,IAAIQ,EAAc,gCAAiC,CAAC,CAAC,CAAC,EAC5GR,EAASC,GAAMM,CAChB,CAEA,IAAI,IAAK,CACR,OAAIP,EAASC,KAAQ,OAAkBD,EAAS,MAAM,IAAIQ,EAAc,oCAAqC,CAAC,CAAC,CAAC,EACzGR,EAASC,EACjB,CAEA,cAAcQ,EAAcC,EAAM,IAAK,CACtC,MAAO,CAAC,KAAK,SAAS,IAAI,KAAMD,CAAI,EAAE,KAAKC,CAAG,CAC/C,CAEA,MAAM,OAAQ,CACb,GAAI,CACH,MAAMC,EAASX,EAASG,GAAO,OAAY,CAAC,CAAC,EAC7C,MAAMQ,EAASX,EAASG,GAAO,OAAY,CAAC,CAAC,CAC9C,OAASS,EAAO,CACfZ,EAAS,MAAM,IAAIQ,EAAc,0BAA2B,CAAC,EAAGI,CAAK,CAAC,CACvE,CACD,CAEA,OAAO,KAAuBC,EAA+B,CAC5D,IAAMC,EAAc,IAAE,SAASD,EAAU,QAAQ,GAAG,EACpD,OAAKC,EAAY,OAChBd,EAAS,MACR,IAAIQ,EAAc;AAAA,EAAwCM,EAAY,MAAM,SAAS,CAAC,GAAI,CACzF,SAAUA,EAAY,MAAM,QAC7B,CAAC,CACF,EAEMA,EAAY,KACpB,CAEA,OAAO,OAAOV,EAAyB,CACtC,GAAIJ,EAASE,GAAW,OAAOF,EAAS,MAAM,IAAIQ,EAAc,wCAAyC,CAAC,CAAC,CAAC,EAC5G,IAAMO,EAAmB,IAAE,SAASC,EAAqB,EAAGZ,CAAQ,EACpE,OAAKW,EAAiB,OACrBf,EAAS,MACR,IAAIQ,EAAc;AAAA,EAA2BO,EAAiB,MAAM,SAAS,CAAC,GAAI,CACjF,SAAUA,EAAiB,MAAM,QAClC,CAAC,CACF,EAEM,IAAIf,EAASe,EAAiB,KAAK,CAC3C,CAEA,OAAO,KAAM,CACZ,OAAKf,EAASE,GAIPF,EAASE,GAHRF,EAAS,MACf,IAAIQ,EAAc,8FAA+F,CAAC,CAAC,CACpH,CAEF,CAEA,OAAO,UAAW,CACjB,OAAOR,EAASE,EACjB,CAEA,OAAO,GAAGe,EAAkBC,EAAYC,EAAuB,CAC9DnB,EAASG,GAAOc,CAAK,IAAM,CAAC,EAC5B,IAAMG,EAAqB,CAAE,GAAAF,EAAI,MAAOC,GAAS,MAAO,MAAOA,GAAS,OAAS,CAAC,CAAE,EACpFE,EAAarB,EAASG,GAAOc,CAAK,EAAGG,CAAM,CAC5C,CAEA,MAAOd,IAAyB,CAO/B,OAAO,QANS,CACf,OAAQ,EACR,OAAQ,EACR,QAAS,EACV,CAEsB,EAAE,QAAQ,CAAC,CAACgB,EAAQC,CAAI,IAAM,CACnD,QAAQ,GAAGD,EAAQ,SAAY,CAC9B,MAAMX,EAASX,EAASG,GAAO,OAAY,CAAC,EAAG,IAAM,CAAC,EAAG,EAAI,EAC7D,QAAQ,KAAK,IAAMoB,CAAI,CACxB,CAAC,CACF,CAAC,CACF,CAEA,OAAO,mBAAsBL,EAAsB,CAClD,IAAMM,EAAQN,EAAG,EACjB,OAAAlB,EAAS,GAAG,QAAS,SAAY,MAAMwB,CAAK,EACrCA,CACR,CAEA,OAAO,MAAMZ,EAA6B,CAEzC,QAAQ,MAAMA,CAAK,EACnB,QAAQ,KAAK,CAAC,CACf,CAEA,OAAO,SAASa,EAAyC,CACxD,MAAO,GAAGA,GAAM,QAAU,EAAE,MAAG,SAAKA,GAAM,MAAM,QAAQ,CAAC,CAAC,EAC3D,CACD,EGnIA,IAAAC,EAA0F,oBAKnF,SAASC,GAA0FC,EAAiBC,EAAY,CACtI,IAAMC,EAAOF,EAAO,EACpB,IAAE,QAAQE,CAAI,EAEd,MAAeC,UAAsBF,CAAgD,CAG1E,YAAgCG,KAA0BC,EAAgC,CAEnG,MAAM,GAAGA,CAAQ,EAFwB,YAAAD,CAG1C,CAEA,OAAO,OAENE,KACGC,EACiB,CACpB,IAAMC,EAAI,IAAE,SAASN,EAAMI,CAAK,EAChC,GAAI,CAACE,EAAE,MAAO,MAAMA,EAAE,MACtB,OAAO,IAAK,KAAaA,EAAE,MAAO,GAAGD,CAAI,CAC1C,CACD,CAEA,OAAOJ,CAQR,CCpCO,IAAMM,GAAkBC,GAAc,CAC5C,GAAI,CACH,OAAIA,GAAM,aAAa,OAAS,SAAiBA,EAC1C,KAAK,MAAMA,CAAI,CACvB,MAAQ,CACP,OAAOA,CACR,CACD,ECPA,IAAAC,EAAA,GAAAC,EAAAD,EAAA,YAAAE,GAAA,WAAAC,KAAA,IAAAC,EAAmB,uBAEZ,SAASD,GAAOE,EAAS,GAAI,CACnC,OAAO,EAAAC,QAAO,YAAYD,CAAM,EAAE,SAAS,KAAK,EAAE,MAAM,EAAGA,CAAM,CAClE,CAEO,SAASH,GAAOK,EAAM,EAAGC,EAAM,GAAK,GAAK,EAAG,CAClD,OAAO,EAAAF,QAAO,UAAUC,EAAKC,CAAG,CACjC,CCNO,IAAMC,GAASC,GAA8B,IAAI,QAASC,GAAY,WAAWA,EAASD,CAAE,CAAC,EAEvFE,EAAQ,MAAUC,EAA+DC,EAAeC,IAAyB,CACrI,GAAID,GAAS,EAAG,MAAM,IAAIE,EAAc,eAAgB,CAAE,MAAAF,EAAO,aAAAC,CAAa,CAAC,EAC/E,IAAME,EAAS,MAAMJ,EAAG,EACxB,OAAII,EAAO,OAAS,GAAaA,EAAO,OACxC,MAAMR,GAAMM,CAAY,EACjB,MAAMH,EAAMC,EAAIC,EAAQ,EAAGC,CAAY,EAC/C,ECVA,IAAAG,GAAkB,sBAClBC,GAAkB,oBCDlB,IAAAC,EAAmC,oBCAnC,IAAAC,EAAwB,0CACxBC,EAAkB,oBAQX,IAAMC,GAAkB,IAC9B,IAAE,KACD,IAAE,OAAO,CACR,QAAS,IAAE,MAAM,IAAE,OAAO,CAAC,EAC3B,IAAK,IAAE,SAAS,IAAE,QAAQ,CAAC,EAC3B,KAAM,IAAE,SACP,IAAE,OAAO,CACR,UAAW,IAAE,GAAG,OAAgB,EAChC,SAAU,IAAE,OAAO,EACnB,SAAU,IAAE,OAAO,CACpB,CAAC,CACF,EACA,SAAU,IAAE,SAAS,IAAE,OAAO,CAAC,CAChC,CAAC,EACD,CAAE,MAAO,eAAgB,OAAQ,aAAc,CAChD,EAEYC,EAAN,MAAMC,UAAsBC,GAAaH,GAAiBI,CAAQ,CAAE,CAC1EC,GACAC,GAAwC,KAE9B,YAAYC,EAAqC,CAC1D,MAAMA,CAAM,EACZ,KAAKF,GAAU,IAAI,UAAQ,MAAM,CAChC,QAAS,CAAE,GAAGE,EAAQ,SAAU,UAAQ,SAAS,OAAQ,CAC1D,CAAC,CACF,CAEA,KAAMC,IAAY,CACjB,OAAK,KAAKF,KACT,KAAKA,IAAU,SAAY,CAC1B,IAAMG,EAAQ,KAAKJ,GAAQ,MAAM,EACjC,aAAMI,EAAM,QAAQ,EACbA,CACR,GAAG,GACG,KAAKH,EACb,CAEA,KAAMI,GAAaC,EAAe,CAEjC,MADc,MAAM,KAAKH,GAAU,GACvB,aAAa,CAAE,OAAQ,CAAC,CAAE,MAAAG,CAAM,CAAC,EAAG,QAAS,GAAK,CAAC,CAChE,CAEA,KAAMC,GAAaC,EAAiB,CAEnC,MADc,MAAM,KAAKL,GAAU,GACvB,aAAa,CAACK,CAAO,CAAC,EAAE,MAAM,IAAM,CAAC,CAAC,CACnD,CAEA,OAA2CC,EAA2BC,EAAkC,CAAC,EAA0B,CAClI,IAAMJ,EAAQI,EAAQ,UAAYD,EAAYE,EAAS,IAAI,EAAE,cAAcF,CAAS,EACpF,MAAO,CACN,QAAS,MAAOG,GAAS,CACxB,IAAMC,EAAW,KAAKb,GAAQ,SAAS,EACvC,aAAMa,EAAS,QAAQ,EACvB,MAAMA,EAAS,KAAK,CACnB,MAAAP,EACA,SAAU,CAAC,CAAE,MAAO,KAAK,UAAUM,CAAI,CAAE,CAAC,CAC3C,CAAC,EACM,EACR,EACA,UAAYE,GAAc,CACzB,IAAMC,EAAY,SAAY,CAC7B,MAAM,KAAKV,GAAaC,CAAK,EAC7B,IAAME,EAAUE,EAAQ,OACrBC,EAAS,IAAI,EAAE,cAAc,GAAGA,EAAS,IAAI,EAAE,EAAE,WAAWK,EAAO,OAAO,EAAE,CAAC,EAAE,EAC/EV,EACGW,EAAW,KAAKjB,GAAQ,SAAS,CAAE,QAAS,CAAE,QAAAQ,CAAQ,CAAE,CAAC,EAE/D,MAAMS,EAAS,QAAQ,EACvB,MAAMA,EAAS,UAAU,CAAE,MAAAX,CAAM,CAAC,EAElC,MAAMW,EAAS,IAAI,CAClB,YAAa,MAAO,CAAE,QAAAC,CAAQ,IAAM,CACnC,MAAMP,EAAS,mBAAmB,SAAY,CACxCO,EAAQ,OACb,MAAMJ,EAAUK,GAAeD,EAAQ,MAAM,SAAS,CAAC,CAAC,CACzD,CAAC,EAAE,MAAOE,GACTT,EAAS,MAAM,IAAIU,EAAc,+BAAgC,CAAE,MAAAf,EAAO,QAAAE,EAAS,QAAAE,CAAQ,EAAGU,CAAK,CAAC,CACrG,CACD,CACD,CAAC,EAEGV,EAAQ,QACXC,EAAS,GAAG,QAAS,SAAY,CAChC,MAAMM,EAAS,WAAW,EAC1B,MAAM,KAAKV,GAAaC,CAAO,CAChC,EAAG,CAAE,MAAOX,CAAc,CAAC,CAC7B,EACAc,EAAS,GAAG,QAASI,EAAW,CAAE,MAAOlB,CAAc,CAAC,CACzD,CACD,CACD,CACD,EDjGO,IAAMyB,GAAqB,IACjC,IAAE,OAAO,CACR,YAAa,IAAE,OAAO,EACtB,SAAU,IAAE,WAAWC,CAA0E,CAClG,CAAC,EDDK,IAAMC,EAAc,aAELC,EAAf,KAA2F,CACjG,YACWC,EACAC,EACAC,EACT,CAHS,YAAAF,EACA,eAAAC,EACA,YAAAC,EAEV,KAAK,OAAS,KAAE,OAAOC,GAAmB,EAAGH,CAAM,CACpD,CAEA,MAAgB,mBAAmBI,EAAaC,EAA8B,CAC7E,IAAMC,EAAW,GAAAC,QAAM,OAAO,CAAE,QAAS,KAAK,OAAO,WAAY,CAAC,EAClE,OAAO,MAAMD,EACX,IAAI,eAAeF,CAAG,UAAW,CACjC,eAAgBN,EAChB,wBAAyB,QACzB,4CAA6C,KAC7C,oCAAqC,KACrC,gBAAiB,8CACjB,+BAAgC,QAChC,kBAAmB,8CACnB,iCAAkC,QAClC,GAAGO,CACJ,CAAC,EACA,KAAK,UACU,MAAMC,EAAS,IAAI,eAAeF,CAAG,SAAS,GAC/C,KAAKA,CAAG,GAAG,QAAQ,WAAWA,CAAG,GAAK,EACpD,EACA,MAAOI,GAAQ,CACf,MAAM,IAAIC,EAAc,8BAA+B,CAAE,IAAAL,CAAI,EAAGI,CAAG,CACpE,CAAC,CACH,CACD,EV7BO,IAAME,EAAN,cAAmGC,CAAwB,CACjIC,GAAW,GAEX,YACCC,EACAC,EACAC,EACAC,EACAC,EACC,CACD,MAAMH,EAAQE,EAAWC,CAAM,EAE/B,IAAMC,EAAWC,GAChBA,EAAK,IACF,CACA,GAAGA,EACH,IAAKA,EAAK,IAAI,MAAWA,EAAK,GAC/B,EACC,OAEEC,EAASL,EAAW,OACpBM,EAAUN,EAAW,eACrBO,EAAY,GAAGF,CAAM,IAAIC,CAAO,GAChCE,EAAQ,GAAGC,CAAW,IAAIF,CAAS,GAEnCG,EAAS,2BACTC,EAAY,CAAE,IAAKD,CAAO,EAEhCX,EAAO,SAAS,OAAOS,EAAgB,CAAE,UAAW,EAAK,CAAC,EAAE,UAAU,MAAOJ,GAA2B,CACvG,IAAMQ,EAAKR,EAAK,GAEZS,EAAS,KAAK,MAAMT,EAAK,QAAU,MAAM,EACzCU,EAAQ,KAAK,MAAMV,EAAK,OAAS,MAAM,EAEvCS,IAAQA,EAASV,EAAQU,CAAM,GAC/BC,IAAOA,EAAQX,EAAQW,CAAK,GAC5B,EAAAD,GAAQ,MAAQH,GAAUI,GAAO,MAAQJ,KAEzCE,IAAO,KAAO,KAAK,UAAU,SAAWE,EAC3C,MAAM,KAAK,UAAU,QAAQ,CAC5B,OAAQ,KACR,MAAO,KAAK,OAAOA,CAAK,CACzB,CAAC,EACOF,IAAO,KAAO,KAAK,UAAU,SAAWC,GAAUC,EAC1D,MAAM,KAAK,UAAU,QAAQ,CAC5B,OAAQ,KAAK,OAAOD,CAAM,EAC1B,MAAO,KAAK,OAAOC,CAAK,EACxB,QAAS,SAAO,KAAK,SAAO,KAAKD,EAAQC,CAAK,CAAC,CAChD,CAAC,EACOF,IAAO,KAAO,KAAK,UAAU,SAAWC,GAChD,MAAM,KAAK,UAAU,QAAQ,CAC5B,OAAQ,KAAK,OAAOA,CAAM,EAC1B,MAAO,IACR,CAAC,EACH,CAAC,EAEDE,EAAS,GAAG,QAAS,SAAY,CAC5B,KAAKlB,KACT,KAAKA,GAAW,GAEhB,MAAMmB,EACL,SACiB,MAAM,KAAK,mBAAmBR,EAAO,CACpD,kBAAmB,iDACnB,eAAgB,4CAChB,4BAA6BV,EAAO,IACpC,0BAA2BS,EAC3B,gBAAiB,cACjB,gBAAiB,aACjB,iBAAkBA,EAClB,wBAAyB,QACzB,qBAAsB,GACvB,CAAC,EAEmB,CAAE,KAAM,GAAM,MAAO,EAAK,GAC9C,MAAMP,EAAW,iBAAiBW,EAAW,CAAE,KAAM,CAAE,QAAAL,CAAQ,CAAS,EAAG,CAAE,OAAQ,EAAK,CAAC,EAC3F,MAAMN,EAAW,iBAAiBW,CAAS,EAC3CI,EAAS,IAAI,EAAE,IAAI,KAAK,8BAA8BR,CAAS,cAAc,EACtE,CAAE,KAAM,EAAM,GAEtB,EACA,GACD,EAAE,MAAOU,GAAQF,EAAS,MAAM,IAAIG,EAAc,6BAA8B,CAAE,UAAAX,CAAU,EAAGU,CAAG,CAAC,CAAC,EACrG,CAAC,CACF,CACD,EahGA,IAAAE,GAAkC,uBAElCC,EASO,mBACPC,GAAkB,oBCZlB,IAAAC,EAAkC,oBAErBC,EAAoB,IAChC,IAAE,KACD,IAAE,OAAO,CACR,IAAK,IAAE,OAAO,CACf,CAAC,EACD,CAAE,MAAO,iBAAkB,OAAQ,eAAgB,CACpD,ECRD,IAAAC,GAA2B,mBCA3B,IAAAC,EAA0F,oBAKnF,IAAKC,OACXA,EAAA,IAAM,MACNA,EAAA,GAAK,KAFMA,OAAA,IAKAC,QACXA,EAAA,GAAK,KACLA,EAAA,IAAM,MACNA,EAAA,GAAK,KACLA,EAAA,IAAM,MACNA,EAAA,GAAK,KACLA,EAAA,GAAK,KACLA,EAAA,GAAK,KACLA,EAAA,IAAM,MARKA,QAAA,IAYNC,GAAY,IAAE,MAAM,IAAE,SAAS,IAAE,GAAG,CAAC,MAAe,IAAY,CAAC,EAAG,KAAa,EAAG,KAAa,EACjGC,GAAa,IAAE,OAAO,CAC3B,MAAO,IAAE,OAAO,EAChB,MAAO,IAAE,IAAI,EAEb,UAAW,IAAE,MAAM,IAAE,SAAS,IAAE,GAAG,OAAO,OAAOF,EAAU,CAAC,EAAG,IAAa,EAAG,IAAa,CAC7F,CAAC,EACKG,GAAkB,IAAE,UACzB,IACC,IAAE,aAAcC,GAAO,OAAO,OAAOL,CAAS,EAAE,SAASK,EAAE,SAAgB,EAAI,QAAU,UAAY,CACpG,MAAO,IAAE,OAAO,CACf,UAAWH,GACX,MAAO,IAAE,MAAME,EAAe,CAC/B,CAAC,EACD,QAASD,EACV,CAAC,EACF,iBACD,EAaMG,GAAmB,IAAE,SAAS,IAAE,MAAMF,EAAe,EAAG,CAAC,CAAC,ED/CzD,IAAMG,GAA0B,MACtCC,EACAC,IACkC,CAElC,IAAMC,EAA8C,CAAC,EAC/CC,EAAQC,EAAgBH,EAAO,MAAOA,EAAO,SAAS,EACxDE,GAAOD,EAAM,KAAKC,CAAK,EAC3B,IAAME,EAAOD,EAAgBH,EAAO,KAAMA,EAAO,QAAQ,EAEzD,GADII,GAAMH,EAAM,KAAKG,CAAI,EACrBJ,EAAO,QAAUA,EAAO,OAAO,OAAO,OAAS,EAAG,CACrD,IAAMK,EAASL,EAAO,OAAO,OAAO,IAAKM,KAAW,CACnD,CAACA,EAAK,EAAG,CACR,OAAQ,IAAI,OAAON,EAAO,OAAQ,MAAO,GAAG,CAC7C,CACD,EAAE,EACFC,EAAM,KAAK,CAAE,IAAKI,CAAO,CAAC,CAC3B,CACA,IAAME,EAAc,CAAC,EACjBN,EAAM,OAAS,IAAGM,EAAY,KAAUN,GAG5C,IAAMO,EAAOR,EAAO,KAAK,IAAKS,GAAM,CAACA,EAAE,MAAOA,EAAE,KAAO,OAAS,KAAK,CAAC,EAGhEC,EAAMV,EAAO,KAAO,GACpBW,EAAQX,EAAO,MACfY,EAAOZ,EAAO,KAEda,EAAQ,MAAMd,EAAW,eAAeQ,CAAW,EAEnDO,EAAad,EAAO,QAAQ,OAAS,OAAO,YAAYA,EAAO,OAAO,IAAKe,GAAM,CAACA,EAAG,CAAC,CAAC,CAAC,EAAI,OAC9FC,EAAajB,EAAW,KAAKQ,EAAa,CAAE,WAAAO,CAAW,CAAC,EACxDN,EAAK,SAAQQ,EAAaA,EAAW,KAAK,OAAO,YAAYR,CAAI,CAAC,GAClE,CAACE,GAAOC,IACXK,EAAaA,EAAW,MAAML,CAAK,EAC/BC,IAAMI,EAAaA,EAAW,MAAMJ,EAAO,GAAKD,CAAK,IAG1D,IAAMM,EAAU,MAAMD,EAAW,QAAQ,EACnCE,EAAQ,EACRC,EAAO,KAAK,KAAKN,EAAQF,CAAK,GAAK,EACnCS,EAAOR,GAAQO,EAAO,KAAOP,EAAO,EACpCS,EAAWT,GAAQM,EAAQ,KAAON,EAAO,EAE/C,MAAO,CACN,MAAO,CAAE,MAAAM,EAAO,KAAAC,EAAM,KAAAC,EAAM,SAAAC,EAAU,QAAST,CAAK,EACpD,KAAM,CAAE,MAAAD,EAAO,MAAAE,EAAO,MAAOI,EAAQ,MAAO,EAC5C,QAASA,CACV,CACD,EAEA,SAASK,GAAaC,EAAmD,CACxE,OAAO,OAAO,OAAOC,CAAS,EAAE,SAASD,EAAM,SAAsB,CACtE,CAEA,IAAMpB,EAAkB,CAACH,EAA4ByB,UAA+E,CACnI,IAAMvB,GAAS,MAAM,QAAQF,CAAM,EAAIA,EAAS,CAAC,GAC/C,IAAKuB,GAAU,CACf,GAAID,GAAaC,CAAK,EAAG,OAAOpB,EAAgBoB,EAAM,MAAOA,EAAM,SAAS,EAC5E,GAAM,CAAE,MAAAjB,CAAM,EAAIiB,EACZG,EAAepB,IAAU,KAAO,MAASA,GAAS,GAClDqB,EAAeJ,EAAM,QAAU,OAAY,GAAKA,EAAM,MAC5D,MAAO,CACN,MAAOG,EACP,MAAOC,EACP,UAAWJ,EAAM,UACjB,QAAS,EACV,CACD,CAAC,EACA,OAAQK,GAAM,CAAC,CAACA,CAAC,EACjB,IAAKA,GACDA,EAAE,QAAgB,CAAE,CAAC,GAAGA,EAAE,KAAK,EAAE,EAAG,CAAE,CAAC,IAAIA,EAAE,SAAS,EAAE,EAAGA,EAAE,KAAM,CAAE,EAC7DA,CACZ,EAEF,OAAO1B,EAAM,OAAS,EAAI,CAAE,CAAC,IAAIuB,CAAG,EAAE,EAAGvB,CAAM,EAAI,IACpD,EE5EO,IAAe2B,EAAf,KAA6C,CACnD,YAAsBC,EAAkB,CAAlB,YAAAA,CAAmB,CAE/B,YAAYC,EAAY,CACjC,OAAOC,EAAS,IAAI,EAAE,cAAcD,CAAE,EAAE,WAAW,IAAK,GAAG,CAC5D,CAOD,EJMA,IAAME,EAAQ,MAGRC,EAAe,IAAI,qBAA6C,MAAS,EAElEC,EAAN,cAAsBC,CAAoB,CAIhD,YACSC,EACRC,EACC,CACD,MAAMA,CAAQ,EAHN,iBAAAD,EAIR,KAAK,YAAc,KAAE,OAAOE,EAAkB,EAAGF,CAAW,EAC5D,KAAK,OAAS,IAAI,cAAYA,EAAY,IAAK,CAAE,gBAAiB,EAAK,CAAC,EACxEG,EAAS,GAAG,QAAS,SAAY,CAChC,MAAM,KAAK,OAAO,QAAQ,EAE1B,IAAMC,EAAU,KAAKC,GAAM,OAAiC,CAACC,EAAKC,KAC5DD,EAAIC,EAAI,EAAE,IAAGD,EAAIC,EAAI,EAAE,EAAI,CAAC,GACjCD,EAAIC,EAAI,EAAE,EAAE,KAAKA,EAAI,GAAG,EACjBD,GACL,CAAC,CAAC,EAECE,EAAU,CACf,6BAA8B,CAAE,QAAS,EAAK,CAC/C,EACA,MAAM,QAAQ,IACb,OAAO,QAAQJ,CAAO,EAAE,IAAI,MAAO,CAACK,EAAQC,CAAQ,IAAM,CACzD,IAAMC,EAAK,KAAK,OAAO,GAAGF,CAAM,EAC1BG,EAAc,MAAMD,EAAG,gBAAgC,EAAE,QAAQ,EACvE,OAAOD,EAAS,IAAI,MAAOG,GAAY,CACtC,IAAMC,EAAWF,EAAY,KAAMG,GAAeA,EAAW,OAASF,CAAO,EACzEC,EAEFA,EAAS,SAAS,8BAA8B,UAAYN,EAAQ,6BAA6B,SAEjG,MAAMG,EAAG,QAAQ,CAAE,QAASE,EAAS,GAAGL,CAAQ,CAAC,EAC5C,MAAMG,EAAG,iBAAiBE,EAASL,CAAO,CAClD,CAAC,CACF,CAAC,CACF,CACD,CAAC,EACDL,EAAS,GAAG,QAAS,SAAY,KAAK,OAAO,MAAM,CAAC,CACrD,CAvCA,OACAE,GAAuC,CAAC,EAwCxC,MAAM,QAAWW,EAA4B,CAC5C,GAAInB,EAAa,SAAS,EAAG,OAAOmB,EAAS,EAC7C,IAAMC,EAAU,MAAM,KAAK,OAAO,aAAa,EAC/C,OAAOA,EAAQ,gBAAgB,SAAYpB,EAAa,IAAIoB,EAASD,CAAQ,CAAC,CAC/E,CAEA,IAAK,CACJ,OAAO,IAAI,UACZ,CAEA,IAA2EE,EAAoC,CAC9G,IAAMP,EAAK,KAAK,YAAYO,EAAO,EAAE,EACrC,YAAKb,GAAM,KAAK,CAAE,GAAAM,EAAI,IAAKO,EAAO,GAAI,CAAC,EAChC,KAAKC,GAAUD,EAAQ,KAAK,OAAO,GAAGP,CAAE,EAAE,WAAkBO,EAAO,GAAG,CAAC,CAC/E,CAEAC,GACCD,EACAH,EACC,CAMD,SAASK,EAAUC,EAAgBC,EAA8B,CAEhE,IAAMC,GADO,MAAM,QAAQF,CAAG,EAAIA,EAAM,CAACA,CAAG,GACxB,IAAKG,GAAMN,EAAO,OAAOM,EAAYF,CAAM,CAAC,EAChE,OAAO,MAAM,QAAQD,CAAG,EAAIE,EAASA,EAAO,CAAC,CAC9C,CAEA,SAASE,EAAgBC,EAAgCC,EAAYC,EAAWC,EAAsB,CACrG,IAAMC,EAA2B,CAChC,CAAClC,CAAK,EAAG+B,EACT,GAAIT,EAAO,SAAS,UACjB,CAAC,EACD,CACA,UAAWU,EAAI,QAAQ,EACvB,GAAIC,EAAa,CAAC,EAAI,CAAE,UAAWD,EAAI,QAAQ,CAAE,CAClD,CACH,EACA,MAAO,CACN,GAAGF,EACH,GAAGI,CACJ,CACD,CAEA,SAASC,EAAgBL,EAAgCE,EAAWI,EAAS,GAAO,CACnF,MAAO,CACN,GAAGN,EACH,KAAM,CACL,GAAGA,EAAM,KACT,GAAIM,GAAW,OAAO,KAAKN,CAAK,EAAE,OAAS,GAAK,CAACR,EAAO,SAAS,UAAa,CAAE,UAAWU,EAAI,QAAQ,CAAE,EAAI,CAAC,CAC/G,CACD,CACD,CAEA,IAAMK,EAAS,KAETC,EAA8E,CACnF,OAAAhB,EACA,OAAQ,CAAE,WAAAH,CAAW,EAErB,MAAO,MAAOoB,GAAW,CACxB,IAAMC,EAAU,MAAMC,GAAwBtB,EAAYoB,CAAyB,EACnF,MAAO,CACN,GAAGC,EACH,QAAShB,EAAUgB,EAAQ,QAAgBD,EAAO,MAA6B,CAChF,CACD,EAEA,SAAU,MAAOG,EAAQ9B,EAAU,CAAC,IAAM,CACzC,GAAM,CAAE,OAAAc,EAAQ,GAAGiB,CAAK,EAAI/B,EACtBgC,EAAalB,GAAQ,OAAS,OAAO,YAAYA,EAAO,IAAKmB,GAAM,CAACA,EAAG,CAAC,CAAC,CAAC,EAAI,OAE9EC,GADY,MAAM,QAAQH,EAAK,IAAI,EAAIA,EAAK,KAAOA,EAAK,KAAO,CAACA,EAAK,IAAI,EAAI,CAAC,GAC7D,IAAKI,GAAM,CAACA,EAAE,MAAOA,EAAE,KAAO,OAAS,KAAK,CAA4B,EACzFC,EAAO,MAAM7B,EACjB,KAAKuB,EAAQ,CACb,QAASzC,EAAa,SAAS,EAC/B,MAAO0C,EAAK,MACZ,KAAAG,EACA,WAAAF,CACD,CAAC,EACA,QAAQ,EACV,OAAOpB,EAAUwB,EAAMtB,CAA6B,CACrD,EAEA,QAAS,MAAOgB,EAA4B9B,EAAU,CAAC,KACvC,MAAM0B,EAAM,SAASI,EAAQ,CAAE,GAAG9B,EAAS,MAAO,CAAE,CAAQ,GAC5D,GAAG,CAAC,GAAK,KAGzB,SAAU,MAAOmB,EAAInB,EAAU,CAAC,IAChB,MAAM0B,EAAM,QAAQ,CAAE,CAACtC,CAAK,EAAG+B,CAAG,EAAyBnB,CAAc,EAIzF,WAAY,MAAOqC,EAAQrC,EAAU,CAAC,IAAM,CAC3C,IAAMoB,EAAMpB,EAAQ,UAAU,GAAK,IAAI,KACjCsC,EAAUD,EAAO,IAAI,CAACnB,EAAOqB,IAAMtB,EAAgBC,EAAOlB,EAAQ,SAASuC,CAAC,GAAK,IAAI,WAAS,EAAE,SAAS,EAAGnB,CAAG,CAAC,EACtH,aAAMb,EAAW,WAAW+B,EAAS,CAAE,QAASjD,EAAa,SAAS,CAAE,CAAC,GAEpD,MAAM,QAAQ,IAAIiD,EAAQ,IAAI,MAAOE,GAAS,MAAMd,EAAM,SAASc,EAAKpD,CAAK,CAAQ,CAAC,CAAC,GACxF,OAAQ8B,GAAU,CAAC,CAACA,CAAK,CAC9C,EAEA,UAAW,MAAOmB,EAAQrC,EAAU,CAAC,KACrB,MAAM0B,EAAM,WAAW,CAACW,CAAM,EAAGrC,CAAO,GACzC,CAAC,EAGhB,WAAY,MAAO8B,EAAQO,EAAQrC,EAAU,CAAC,IAAM,CACnD,IAAMoB,EAAMpB,EAAQ,UAAU,GAAK,IAAI,KACjCS,EAAUpB,EAAa,SAAS,EAEhCoD,GADO,MAAMlC,EAAW,KAAKuB,EAAQ,CAAE,QAAArB,EAAS,WAAY,CAAE,CAACrB,CAAK,EAAG,CAAE,CAAE,CAAC,EAAE,QAAQ,GAC3E,IAAKyB,GAAQA,EAAIzB,CAAK,CAAC,EAClCsD,EAAY,CAAE,CAACtD,CAAK,EAAG,CAAE,IAAKqD,CAAI,CAAE,EAC1C,aAAMlC,EAAW,WAAWmC,EAAWnB,EAAgBc,EAAQjB,CAAG,EAAG,CAAE,QAAAX,CAAQ,CAAC,EACzEiB,EAAM,SAASgB,CAAS,CAChC,EAEA,UAAW,MAAOZ,EAAQO,EAAQrC,EAAU,CAAC,IAAM,CAClD,IAAMoB,EAAMpB,EAAQ,UAAU,GAAK,IAAI,KACjCa,EAAM,MAAMN,EAAW,iBAAiBuB,EAAQP,EAAgBc,EAAQjB,CAAG,EAAG,CACnF,eAAgB,QAChB,QAAS/B,EAAa,SAAS,CAChC,CAAC,EACD,OAAOwB,EAAMD,EAAUC,EAAK,MAAS,EAAI,IAC1C,EAEA,WAAY,MAAOM,EAAIkB,EAAQrC,EAAU,CAAC,IAC1B,MAAM0B,EAAM,UAAU,CAAE,CAACtC,CAAK,EAAG+B,CAAG,EAAyBkB,EAAQrC,CAAO,EAI5F,UAAW,MAAO8B,EAAQO,EAAQrC,EAAU,CAAC,IAAM,CAClD,IAAMoB,EAAMpB,EAAQ,UAAU,GAAK,IAAI,KAEjCa,EAAM,MAAMN,EAAW,iBAC5BuB,EACA,CACC,GAAGP,EAAgB,WAAYc,EAASA,EAAO,OAAS,CAAC,EAAGjB,EAAK,EAAI,EAErE,aAAcH,EAAgBoB,EAAO,OAAQrC,EAAQ,SAAS,GAAK,IAAI,WAAS,EAAE,SAAS,EAAGoB,EAAK,EAAI,CACxG,EACA,CAAE,eAAgB,QAAS,QAAS/B,EAAa,SAAS,EAAG,OAAQ,EAAK,CAC3E,EAEA,OAAOuB,EAAUC,EAAK,MAAS,CAChC,EAEA,WAAY,MAAOiB,EAAQ9B,EAAU,CAAC,IAAM,CAC3C,IAAMoC,EAAO,MAAMV,EAAM,SAASI,EAAQ9B,CAAO,EACjD,aAAMO,EAAW,WAAWuB,EAAQ,CAAE,QAASzC,EAAa,SAAS,CAAE,CAAC,EACjE+C,CACR,EAEA,UAAW,MAAON,GAAW,CAC5B,IAAMjB,EAAM,MAAMN,EAAW,iBAAiBuB,EAAQ,CAAE,QAASzC,EAAa,SAAS,CAAE,CAAC,EAC1F,OAAOwB,EAAMD,EAAUC,EAAK,MAAS,EAAI,IAC1C,EAEA,WAAY,MAAOM,GACH,MAAMO,EAAM,UAAU,CAAE,CAACtC,CAAK,EAAG+B,CAAG,CAAuB,EAI3E,UAAW,MAAOwB,EAAY3C,EAAU,CAAC,IAAM,CAC9C,GAAI,CAAC2C,EAAW,OAAQ,OACxB,IAAMC,EAAOrC,EAAW,0BAA0B,CAAE,QAASlB,EAAa,SAAS,CAAE,CAAC,EAChF+B,EAAMpB,EAAQ,UAAU,GAAK,IAAI,KACvC2C,EAAW,QAAQ,CAACE,EAAWN,IAAM,CACpC,OAAQM,EAAU,GAAI,CACrB,IAAK,SACJD,EAAK,OAAO3B,EAAgB4B,EAAU,MAAOA,EAAU,SAASN,CAAC,GAAK,IAAI,WAAS,EAAE,SAAS,EAAGnB,CAAG,CAAC,EACrG,MACD,IAAK,SACJwB,EAAK,KAAKC,EAAU,MAAM,EAAE,OAAO,EACnC,MACD,IAAK,SACJD,EAAK,KAAKC,EAAU,MAAM,EAAE,OAAOtB,EAAgBsB,EAAU,MAAOzB,CAAG,CAAC,EACxE,MACD,IAAK,SACJwB,EAAK,KAAKC,EAAU,MAAM,EACxB,OAAO,EACP,OAAO,CACP,GAAGtB,EAAgB,WAAYsB,EAAYA,EAAU,OAAS,CAAC,EAAGzB,EAAK,EAAI,EAC3E,aAAcH,EACb4B,EAAU,OACVA,EAAU,SAASN,CAAC,GAAK,IAAI,WAAS,EAAE,SAAS,EACjDnB,EACA,EACD,CACD,CAAC,EACF,MACD,QACC,MAAM,IAAI0B,EAAc,8BAA+B,CAAE,UAAAD,CAAU,CAAC,CACtE,CACD,CAAC,EACD,MAAMD,EAAK,QAAQ,CAAE,QAASvD,EAAa,SAAS,CAAE,CAAC,CACxD,EAEA,MAAM0D,EAAW,CAChB,OAAKtB,EAAO,OAAO,SAClB9B,EAAS,MAAM,IAAImD,EAAc,mDAAoD,CAAE,OAAApC,CAAO,CAAC,CAAC,EAC1F,IAAIsC,EAA6BvB,EAAO,YAAaA,EAAO,OAAO,QAASlB,EAAYwC,EAAYE,GAC1GvC,EAAO,OAAOuC,CAAC,CAChB,CACD,CACD,EAEA,OAAOvB,CACR,CACD","names":["mongodb_exports","__export","MongoDb","MongoDbChange","mongoDbConfigPipe","__toCommonJS","import_mongodb","import_valleyed","EquippedError","message","context","cause","EventBus","import_pino","import_ulid","import_valleyed","registerHook","hooks","record","dep","h","depName","ownerName","tentative","detectCycle","classToAfter","existing","visited","stack","visit","cls","path","cycleStart","names","c","resolveHookDAG","invert","classToHooks","anonymous","list","classes","inDegree","graph","from","to","edges","layers","queue","layer","next","neighbor","d","anonWithDeps","anonWithoutDeps","maxLayer","depLayer","l","r","insertAt","runHooks","onError","error","import_valleyed","instanceSettingsPipe","Instance","_Instance","#id","#instance","#hooks","settings","pino","#registerOnExitHandler","id","EquippedError","name","key","runHooks","error","envsPipe","envValidity","settingsValidity","instanceSettingsPipe","event","cb","options","record","registerHook","signal","code","value","opts","import_valleyed","configurable","pipeFn","base","pipe","Configurable","config","baseArgs","input","args","r","parseJSONValue","data","random_exports","__export","number","string","import_crypto","length","crypto","min","max","sleep","ms","resolve","retry","cb","tries","waitTimeInMs","EquippedError","result","import_axios","import_valleyed","import_valleyed","import_kafka_javascript","import_valleyed","kafkaConfigPipe","KafkaEventBus","_KafkaEventBus","configurable","EventBus","#client","#admin","config","#getAdmin","admin","#createTopic","topic","#deleteGroup","groupId","topicName","options","Instance","data","producer","onMessage","subscribe","random_exports","consumer","message","parseJSONValue","error","EquippedError","dbChangeConfigPipe","KafkaEventBus","TopicPrefix","DbChange","config","callbacks","mapper","dbChangeConfigPipe","key","data","instance","axios","err","EquippedError","MongoDbChange","DbChange","#started","config","change","collection","callbacks","mapper","hydrate","data","dbName","colName","dbColName","topic","TopicPrefix","TestId","condition","op","before","after","Instance","retry","err","EquippedError","import_node_async_hooks","import_mongodb","import_valleyed","import_valleyed","mongoDbConfigPipe","import_mongodb","import_valleyed","QueryKeys","Conditions","queryKeys","queryWhere","queryWhereBlock","d","queryWhereClause","parseMongodbQueryParams","collection","params","query","where","buildWhereQuery","auth","search","field","totalClause","sort","p","all","limit","page","total","projection","k","builtQuery","results","start","last","next","previous","isWhereBlock","param","QueryKeys","key","checkedField","checkedValue","c","Db","config","db","Instance","idKey","sessionStore","MongoDb","Db","mongoConfig","dbConfig","mongoDbConfigPipe","Instance","grouped","#cols","acc","cur","options","dbName","colNames","db","collections","colName","existing","collection","callback","session","config","#getTable","transform","doc","select","mapped","d","prepInsertValue","value","id","now","skipUpdate","base","prepUpdateValue","upsert","dbThis","table","params","results","parseMongodbQueryParams","filter","rest","projection","k","sort","p","docs","values","payload","i","data","ids","filterUpd","operations","bulk","operation","EquippedError","callbacks","MongoDbChange","m"]}