diff --git a/src/config/config.ts b/src/config/config.ts index 9f5e922c..d14e61ce 100644 --- a/src/config/config.ts +++ b/src/config/config.ts @@ -59,8 +59,8 @@ export function getDefaultConfig(): PluginsServerConfig { DISTINCT_ID_LRU_SIZE: 10000, INTERNAL_MMDB_SERVER_PORT: 0, PLUGIN_SERVER_IDLE: false, - RETRY_QUEUES: '', - RETRY_QUEUE_GRAPHILE_URL: '', + JOB_QUEUES: '', + JOB_QUEUE_GRAPHILE_URL: '', ENABLE_PERSISTENT_CONSOLE: false, // TODO: remove when persistent console ships in main repo STALENESS_RESTART_SECONDS: 0, } @@ -104,8 +104,8 @@ export function getConfigHelp(): Record { DISTINCT_ID_LRU_SIZE: 'size of persons distinct ID LRU cache', INTERNAL_MMDB_SERVER_PORT: 'port of the internal server used for IP location (0 means random)', PLUGIN_SERVER_IDLE: 'whether to disengage the plugin server, e.g. for development', - RETRY_QUEUES: 'retry queue engine and fallback queues', - RETRY_QUEUE_GRAPHILE_URL: 'use a different postgres connection in the graphile retry queue', + JOB_QUEUES: 'retry queue engine and fallback queues', + JOB_QUEUE_GRAPHILE_URL: 'use a different postgres connection in the graphile retry queue', STALENESS_RESTART_SECONDS: 'trigger a restart if no event ingested for this duration', } } diff --git a/src/main/job-queues/fs-queue.ts b/src/main/job-queues/fs-queue.ts index 97714f48..c1d3fb1c 100644 --- a/src/main/job-queues/fs-queue.ts +++ b/src/main/job-queues/fs-queue.ts @@ -1,4 +1,4 @@ -import { EnqueuedRetry, JobQueue, OnRetryCallback } from '../../types' +import { EnqueuedJob, JobQueue, OnJobCallback } from '../../types' import Timeout = NodeJS.Timeout import * as fs from 'fs' import * as path from 'path' @@ -17,20 +17,22 @@ export class FsQueue implements JobQueue { this.started = false this.interval = null this.filename = filename || path.join(process.cwd(), 'tmp', 'fs-queue.txt') + } + connectProducer(): void { fs.mkdirSync(path.dirname(this.filename), { recursive: true }) fs.writeFileSync(this.filename, '') } - enqueue(retry: EnqueuedRetry): Promise | void { - fs.appendFileSync(this.filename, `${JSON.stringify(retry)}\n`) + enqueue(job: EnqueuedJob): Promise | void { + fs.appendFileSync(this.filename, `${JSON.stringify(job)}\n`) } - quit(): void { + disconnectProducer(): void { // nothing to do } - startConsumer(onRetry: OnRetryCallback): void { + startConsumer(onJob: OnJobCallback): void { fs.writeFileSync(this.filename, '') this.started = true this.interval = setInterval(() => { @@ -43,14 +45,14 @@ export class FsQueue implements JobQueue { .toString() .split('\n') .filter((a) => a) - .map((s) => JSON.parse(s) as EnqueuedRetry) + .map((s) => JSON.parse(s) as EnqueuedJob) const newQueue = queue.filter((element) => element.timestamp < timestamp) if (newQueue.length > 0) { const oldQueue = queue.filter((element) => element.timestamp >= timestamp) fs.writeFileSync(this.filename, `${oldQueue.map((q) => JSON.stringify(q)).join('\n')}\n`) - void onRetry(newQueue) + void onJob(newQueue) } }, 1000) } diff --git a/src/main/job-queues/graphile-queue.ts b/src/main/job-queues/graphile-queue.ts index 97c3b298..42cb1ccd 100644 --- a/src/main/job-queues/graphile-queue.ts +++ b/src/main/job-queues/graphile-queue.ts @@ -1,49 +1,44 @@ import { makeWorkerUtils, run, Runner, WorkerUtils, WorkerUtilsOptions } from 'graphile-worker' -import { EnqueuedRetry, JobQueue, OnRetryCallback, PluginsServer } from '../../types' +import { EnqueuedJob, JobQueue, OnJobCallback, PluginsServer } from '../../types' export class GraphileQueue implements JobQueue { pluginsServer: PluginsServer started: boolean paused: boolean - onRetry: OnRetryCallback | null + onJob: OnJobCallback | null runner: Runner | null - workerUtils: WorkerUtils | null + workerUtilsPromise: Promise | null constructor(pluginsServer: PluginsServer) { this.pluginsServer = pluginsServer this.started = false this.paused = false - this.onRetry = null + this.onJob = null this.runner = null - this.workerUtils = null + this.workerUtilsPromise = null } - async enqueue(retry: EnqueuedRetry): Promise { - if (!this.workerUtils) { - this.workerUtils = await makeWorkerUtils( - this.pluginsServer.RETRY_QUEUE_GRAPHILE_URL - ? { - connectionString: this.pluginsServer.RETRY_QUEUE_GRAPHILE_URL, - } - : ({ - pgPool: this.pluginsServer.postgres, - } as WorkerUtilsOptions) - ) - await this.workerUtils.migrate() - } - await this.workerUtils.addJob('retryTask', retry, { runAt: new Date(retry.timestamp), maxAttempts: 1 }) + async connectProducer(): Promise { + await (await this.getWorkerUtils()).migrate() + } + + async enqueue(retry: EnqueuedJob): Promise { + await (await this.getWorkerUtils()).addJob('pluginJob', retry, { + runAt: new Date(retry.timestamp), + maxAttempts: 1, + }) } - async quit(): Promise { - const oldWorkerUtils = this.workerUtils - this.workerUtils = null + async disconnectProducer(): Promise { + const oldWorkerUtils = await this.workerUtilsPromise + this.workerUtilsPromise = null await oldWorkerUtils?.release() } - async startConsumer(onRetry: OnRetryCallback): Promise { + async startConsumer(onJob: OnJobCallback): Promise { this.started = true - this.onRetry = onRetry + this.onJob = onJob await this.syncState() } @@ -66,19 +61,19 @@ export class GraphileQueue implements JobQueue { await this.syncState() } - async syncState(): Promise { + private async syncState(): Promise { if (this.started && !this.paused) { if (!this.runner) { this.runner = await run({ - connectionString: this.pluginsServer.DATABASE_URL, + ...this.getConnectionOptions(), concurrency: 1, // Install signal handlers for graceful shutdown on SIGINT, SIGTERM, etc noHandleSignals: false, pollInterval: 100, // you can set the taskList or taskDirectory but not both taskList: { - retryTask: (payload) => { - void this.onRetry?.([payload as EnqueuedRetry]) + pluginJob: (payload) => { + void this.onJob?.([payload as EnqueuedJob]) }, }, }) @@ -91,4 +86,21 @@ export class GraphileQueue implements JobQueue { } } } + + private getConnectionOptions(): Partial { + return this.pluginsServer.JOB_QUEUE_GRAPHILE_URL + ? { + connectionString: this.pluginsServer.JOB_QUEUE_GRAPHILE_URL, + } + : ({ + pgPool: this.pluginsServer.postgres, + } as Partial) + } + + private async getWorkerUtils(): Promise { + if (!this.workerUtilsPromise) { + this.workerUtilsPromise = makeWorkerUtils(this.getConnectionOptions()) + } + return await this.workerUtilsPromise + } } diff --git a/src/main/job-queues/job-queue-consumer.ts b/src/main/job-queues/job-queue-consumer.ts index 0d02cb97..6943ad5e 100644 --- a/src/main/job-queues/job-queue-consumer.ts +++ b/src/main/job-queues/job-queue-consumer.ts @@ -1,19 +1,19 @@ import Piscina from '@posthog/piscina' -import { JobQueueConsumerControl, OnRetryCallback, PluginsServer } from '../../types' +import { JobQueueConsumerControl, OnJobCallback, PluginsServer } from '../../types' import { startRedlock } from '../../utils/redlock' import { status } from '../../utils/status' import { pauseQueueIfWorkerFull } from '../ingestion-queues/queue' -export const LOCKED_RESOURCE = 'plugin-server:locks:retry-queue-consumer' +export const LOCKED_RESOURCE = 'plugin-server:locks:job-queue-consumer' export async function startJobQueueConsumer(server: PluginsServer, piscina: Piscina): Promise { - status.info('๐Ÿ”„', 'Starting retry queue consumer, trying to get lock...') + status.info('๐Ÿ”„', 'Starting job queue consumer, trying to get lock...') - const onRetry: OnRetryCallback = async (retries) => { - pauseQueueIfWorkerFull(server.retryQueueManager.pauseConsumer, server, piscina) - for (const retry of retries) { - await piscina.runTask({ task: 'retry', args: { retry } }) + const onJob: OnJobCallback = async (jobs) => { + pauseQueueIfWorkerFull(server.jobQueueManager.pauseConsumer, server, piscina) + for (const job of jobs) { + await piscina.runTask({ task: 'runJob', args: { job } }) } } @@ -21,15 +21,15 @@ export async function startJobQueueConsumer(server: PluginsServer, piscina: Pisc server, resource: LOCKED_RESOURCE, onLock: async () => { - status.info('๐Ÿ”„', 'Retry queue consumer lock aquired') - await server.retryQueueManager.startConsumer(onRetry) + status.info('๐Ÿ”„', 'Job queue consumer lock aquired') + await server.jobQueueManager.startConsumer(onJob) }, onUnlock: async () => { - status.info('๐Ÿ”„', 'Stopping retry queue consumer') - await server.retryQueueManager.stopConsumer() + status.info('๐Ÿ”„', 'Stopping job queue consumer') + await server.jobQueueManager.stopConsumer() }, ttl: server.SCHEDULE_LOCK_TTL, }) - return { stop: () => unlock(), resume: () => server.retryQueueManager.resumeConsumer() } + return { stop: () => unlock(), resume: () => server.jobQueueManager.resumeConsumer() } } diff --git a/src/main/job-queues/job-queue-manager.ts b/src/main/job-queues/job-queue-manager.ts index 974a7794..f6929773 100644 --- a/src/main/job-queues/job-queue-manager.ts +++ b/src/main/job-queues/job-queue-manager.ts @@ -1,6 +1,7 @@ import * as Sentry from '@sentry/node' -import { EnqueuedRetry, JobQueue, OnRetryCallback, PluginsServer } from '../../types' +import { EnqueuedJob, JobQueue, OnJobCallback, PluginsServer } from '../../types' +import { status } from '../../utils/status' import { FsQueue } from './fs-queue' import { GraphileQueue } from './graphile-queue' @@ -17,35 +18,50 @@ const queues: Record JobQueue> = { export class JobQueueManager implements JobQueue { pluginsServer: PluginsServer jobQueues: JobQueue[] + jobQueueTypes: JobQueueType[] constructor(pluginsServer: PluginsServer) { this.pluginsServer = pluginsServer - this.jobQueues = pluginsServer.RETRY_QUEUES.split(',') + this.jobQueueTypes = pluginsServer.JOB_QUEUES.split(',') .map((q) => q.trim() as JobQueueType) .filter((q) => !!q) - .map( - (queue): JobQueue => { - if (queues[queue]) { - return queues[queue](pluginsServer) - } else { - throw new Error(`Unknown retry queue "${queue}"`) - } + + this.jobQueues = this.jobQueueTypes.map( + (queue): JobQueue => { + if (queues[queue]) { + return queues[queue](pluginsServer) + } else { + throw new Error(`Unknown job queue "${queue}"`) + } + } + ) + } + + async connectProducer(): Promise { + await Promise.all( + this.jobQueues.map(async (jobQueue, index) => { + try { + await jobQueue.connectProducer() + status.info('๐Ÿ’‚', `Connected to job queue producer: ${this.jobQueueTypes[index]}`) + } catch (error) { + Sentry.captureException(error) } - ) + }) + ) } - async enqueue(retry: EnqueuedRetry): Promise { - for (const retryQueue of this.jobQueues) { + async enqueue(job: EnqueuedJob): Promise { + for (const jobQueue of this.jobQueues) { try { - await retryQueue.enqueue(retry) + await jobQueue.enqueue(job) return } catch (error) { // if one fails, take the next queue Sentry.captureException(error, { extra: { - retry: JSON.stringify(retry), - queue: retryQueue.toString(), + job: JSON.stringify(job), + queue: jobQueue.toString(), queues: this.jobQueues.map((q) => q.toString()), }, }) @@ -54,12 +70,12 @@ export class JobQueueManager implements JobQueue { throw new Error('No JobQueue available') } - async quit(): Promise { - await Promise.all(this.jobQueues.map((r) => r.quit())) + async disconnectProducer(): Promise { + await Promise.all(this.jobQueues.map((r) => r.disconnectProducer())) } - async startConsumer(onRetry: OnRetryCallback): Promise { - await Promise.all(this.jobQueues.map((r) => r.startConsumer(onRetry))) + async startConsumer(onJob: OnJobCallback): Promise { + await Promise.all(this.jobQueues.map((r) => r.startConsumer(onJob))) } async stopConsumer(): Promise { diff --git a/src/main/services/redlock.ts b/src/main/services/redlock.ts deleted file mode 100644 index 03cd2dd9..00000000 --- a/src/main/services/redlock.ts +++ /dev/null @@ -1,100 +0,0 @@ -import * as Sentry from '@sentry/node' -import Redlock from 'redlock' - -import { PluginsServer } from '../../types' -import { status } from '../../utils/status' -import { createRedis } from '../../utils/utils' - -type RedlockOptions = { - server: PluginsServer - resource: string - onLock: () => Promise | void - onUnlock: () => Promise | void - ttl: number -} - -export async function startRedlock({ - server, - resource, - onLock, - onUnlock, - ttl, -}: RedlockOptions): Promise<() => Promise> { - status.info('โฐ', `Starting redlock "${resource}" ...`) - - let stopped = false - let weHaveTheLock = false - let lock: Redlock.Lock - let lockTimeout: NodeJS.Timeout - - const lockTTL = ttl * 1000 // 60 sec if default passed in - const retryDelay = lockTTL / 10 // 6 sec - const extendDelay = lockTTL / 2 // 30 sec - - // use another redis connection for redlock - const redis = await createRedis(server) - - const redlock = new Redlock([redis], { - // we handle retries ourselves to have a way to cancel the promises on quit - // without this, the `await redlock.lock()` code will remain inflight and cause issues - retryCount: 0, - }) - - redlock.on('clientError', (error) => { - if (stopped) { - return - } - status.error('๐Ÿ”ด', `Redlock "${resource}" client error occurred:\n`, error) - Sentry.captureException(error, { extra: { resource } }) - }) - - const tryToGetTheLock = async () => { - try { - lock = await redlock.lock(resource, lockTTL) - weHaveTheLock = true - - status.info('๐Ÿ”’', `Redlock "${resource}" acquired!`) - - const extendLock = async () => { - if (stopped) { - return - } - try { - lock = await lock.extend(lockTTL) - lockTimeout = setTimeout(extendLock, extendDelay) - } catch (error) { - status.error('๐Ÿ”ด', `Redlock cannot extend lock "${resource}":\n`, error) - Sentry.captureException(error, { extra: { resource } }) - weHaveTheLock = false - lockTimeout = setTimeout(tryToGetTheLock, 0) - } - } - - lockTimeout = setTimeout(extendLock, extendDelay) - - await onLock?.() - } catch (error) { - if (stopped) { - return - } - weHaveTheLock = false - if (error instanceof Redlock.LockError) { - lockTimeout = setTimeout(tryToGetTheLock, retryDelay) - } else { - Sentry.captureException(error, { extra: { resource } }) - status.error('๐Ÿ”ด', `Redlock "${resource}" error:\n`, error) - } - } - } - - lockTimeout = setTimeout(tryToGetTheLock, 0) - - return async () => { - stopped = true - lockTimeout && clearTimeout(lockTimeout) - - await lock?.unlock().catch(Sentry.captureException) - await redis.quit() - await onUnlock?.() - } -} diff --git a/src/main/services/retry-queue-consumer.ts b/src/main/services/retry-queue-consumer.ts deleted file mode 100644 index 2dc96abc..00000000 --- a/src/main/services/retry-queue-consumer.ts +++ /dev/null @@ -1,38 +0,0 @@ -import Piscina from '@posthog/piscina' - -import { JobQueueConsumerControl, OnRetryCallback, PluginsServer } from '../../types' -import { status } from '../../utils/status' -import { pauseQueueIfWorkerFull } from '../ingestion-queues/queue' -import { startRedlock } from './redlock' - -export const LOCKED_RESOURCE = 'plugin-server:locks:retry-queue-consumer' - -export async function startRetryQueueConsumer( - server: PluginsServer, - piscina: Piscina -): Promise { - status.info('๐Ÿ”„', 'Starting retry queue consumer, trying to get lock...') - - const onRetry: OnRetryCallback = async (retries) => { - pauseQueueIfWorkerFull(server.retryQueueManager.pauseConsumer, server, piscina) - for (const retry of retries) { - await piscina.runTask({ task: 'retry', args: { retry } }) - } - } - - const unlock = await startRedlock({ - server, - resource: LOCKED_RESOURCE, - onLock: async () => { - status.info('๐Ÿ”„', 'Retry queue consumer lock aquired') - await server.retryQueueManager.startConsumer(onRetry) - }, - onUnlock: async () => { - status.info('๐Ÿ”„', 'Stopping retry queue consumer') - await server.retryQueueManager.stopConsumer() - }, - ttl: server.SCHEDULE_LOCK_TTL, - }) - - return { stop: () => unlock(), resume: () => server.retryQueueManager.resumeConsumer() } -} diff --git a/src/main/services/schedule.ts b/src/main/services/schedule.ts index 32853708..d3546205 100644 --- a/src/main/services/schedule.ts +++ b/src/main/services/schedule.ts @@ -25,19 +25,19 @@ export async function startSchedule( !stopped && weHaveTheLock && (await pluginSchedulePromise) && - runTasksDebounced(server!, piscina!, 'runEveryMinute') + runScheduleDebounced(server!, piscina!, 'runEveryMinute') }) const runEveryHourJob = schedule.scheduleJob('0 * * * *', async () => { !stopped && weHaveTheLock && (await pluginSchedulePromise) && - runTasksDebounced(server!, piscina!, 'runEveryHour') + runScheduleDebounced(server!, piscina!, 'runEveryHour') }) const runEveryDayJob = schedule.scheduleJob('0 0 * * *', async () => { !stopped && weHaveTheLock && (await pluginSchedulePromise) && - runTasksDebounced(server!, piscina!, 'runEveryDay') + runScheduleDebounced(server!, piscina!, 'runEveryDay') }) const unlock = await startRedlock({ @@ -91,7 +91,7 @@ export async function loadPluginSchedule( throw new Error('Could not load plugin schedule in time') } -export function runTasksDebounced(server: PluginsServer, piscina: Piscina, taskName: string): void { +export function runScheduleDebounced(server: PluginsServer, piscina: Piscina, taskName: string): void { const runTask = (pluginConfigId: PluginConfigId) => piscina.runTask({ task: taskName, args: { pluginConfigId } }) for (const pluginConfigId of server.pluginSchedule?.[taskName] || []) { diff --git a/src/types.ts b/src/types.ts index 355c0f8a..39d9c90b 100644 --- a/src/types.ts +++ b/src/types.ts @@ -71,8 +71,8 @@ export interface PluginsServerConfig extends Record { DISTINCT_ID_LRU_SIZE: number INTERNAL_MMDB_SERVER_PORT: number PLUGIN_SERVER_IDLE: boolean - RETRY_QUEUES: string - RETRY_QUEUE_GRAPHILE_URL: string + JOB_QUEUES: string + JOB_QUEUE_GRAPHILE_URL: string ENABLE_PERSISTENT_CONSOLE: boolean STALENESS_RESTART_SECONDS: number } @@ -115,8 +115,8 @@ export interface Queue extends Pausable { stop: () => Promise | void } -export type OnRetryCallback = (queue: EnqueuedRetry[]) => Promise | void -export interface EnqueuedRetry { +export type OnJobCallback = (queue: EnqueuedJob[]) => Promise | void +export interface EnqueuedJob { type: string payload: Record timestamp: number @@ -125,13 +125,15 @@ export interface EnqueuedRetry { } export interface JobQueue { - startConsumer: (onRetry: OnRetryCallback) => Promise | void + startConsumer: (onJob: OnJobCallback) => Promise | void stopConsumer: () => Promise | void pauseConsumer: () => Promise | void resumeConsumer: () => Promise | void isConsumerPaused: () => boolean - enqueue: (retry: EnqueuedRetry) => Promise | void - quit: () => Promise | void + + connectProducer: () => Promise | void + enqueue: (job: EnqueuedJob) => Promise | void + disconnectProducer: () => Promise | void } export type PluginId = number @@ -226,10 +228,15 @@ export interface PluginLogEntry { instance_id: string } +export enum PluginTaskType { + Job = 'job', + Schedule = 'schedule', +} + export interface PluginTask { name: string - type: 'runEvery' - exec: () => Promise + type: PluginTaskType + exec: (payload?: Record) => Promise } export interface PluginConfigVMReponse { @@ -239,9 +246,8 @@ export interface PluginConfigVMReponse { teardownPlugin: () => Promise processEvent: (event: PluginEvent) => Promise processEventBatch: (batch: PluginEvent[]) => Promise - onRetry: (task: string, payload: Record) => Promise } - tasks: Record + tasks: Record> } export interface EventUsage { diff --git a/src/utils/db/server.ts b/src/utils/db/server.ts index e356f531..fd1803fb 100644 --- a/src/utils/db/server.ts +++ b/src/utils/db/server.ts @@ -167,11 +167,12 @@ export async function createServer( // :TODO: This is only used on worker threads, not main server.eventsProcessor = new EventsProcessor(server as PluginsServer) - server.retryQueueManager = new JobQueueManager(server as PluginsServer) + server.jobQueueManager = new JobQueueManager(server as PluginsServer) + await server.jobQueueManager.connectProducer() const closeServer = async () => { server.mmdbUpdateJob?.cancel() - await server.retryQueueManager?.quit() + await server.jobQueueManager?.disconnectProducer() if (kafkaProducer) { clearInterval(kafkaProducer.flushInterval) await kafkaProducer.flush() diff --git a/src/utils/redlock.ts b/src/utils/redlock.ts index 17643371..e4145ae8 100644 --- a/src/utils/redlock.ts +++ b/src/utils/redlock.ts @@ -90,6 +90,9 @@ export async function startRedlock({ lockTimeout = setTimeout(tryToGetTheLock, 0) return async () => { + if (weHaveTheLock) { + status.info('๐Ÿ”“', `Releasing redlock "${resource}"`) + } stopped = true lockTimeout && clearTimeout(lockTimeout) diff --git a/src/worker/plugins/run.ts b/src/worker/plugins/run.ts index 613f3666..59b19b04 100644 --- a/src/worker/plugins/run.ts +++ b/src/worker/plugins/run.ts @@ -1,7 +1,6 @@ import { PluginEvent } from '@posthog/plugin-scaffold' -import * as Sentry from '@sentry/node' -import { EnqueuedRetry, PluginConfig, PluginsServer } from '../../types' +import { PluginConfig, PluginsServer, PluginTaskType } from '../../types' import { processError } from '../../utils/db/error' export async function runPlugins(server: PluginsServer, event: PluginEvent): Promise { @@ -98,43 +97,32 @@ export async function runPluginsOnBatch(server: PluginsServer, batch: PluginEven return allReturnedEvents.filter(Boolean) } -export async function runPluginTask(server: PluginsServer, taskName: string, pluginConfigId: number): Promise { +export async function runPluginTask( + server: PluginsServer, + taskName: string, + taskType: PluginTaskType, + pluginConfigId: number, + payload?: Record +): Promise { const timer = new Date() let response const pluginConfig = server.pluginConfigs.get(pluginConfigId) try { - const task = await pluginConfig?.vm?.getTask(taskName) - response = await task?.exec() + const task = await pluginConfig?.vm?.getTask(taskName, taskType) + if (!task) { + throw new Error( + `Task "${taskName}" not found for plugin "${pluginConfig?.plugin?.name}" with config id ${pluginConfig}` + ) + } + response = await (payload ? task?.exec(payload) : task?.exec()) } catch (error) { await processError(server, pluginConfig || null, error) - server.statsd?.increment(`plugin.task.${taskName}.${pluginConfigId}.ERROR`) + server.statsd?.increment(`plugin.task.${taskType}.${taskName}.${pluginConfigId}.ERROR`) } - server.statsd?.timing(`plugin.task.${taskName}.${pluginConfigId}`, timer) + server.statsd?.timing(`plugin.task.${taskType}.${taskName}.${pluginConfigId}`, timer) return response } function getPluginsForTeam(server: PluginsServer, teamId: number): PluginConfig[] { return server.pluginConfigsPerTeam.get(teamId) || [] } - -export async function runOnRetry(server: PluginsServer, retry: EnqueuedRetry): Promise { - const timer = new Date() - let response - const pluginConfig = server.pluginConfigs.get(retry.pluginConfigId) - if (pluginConfig) { - try { - const task = await pluginConfig.vm?.getOnRetry() - response = await task?.(retry.type, retry.payload) - } catch (error) { - await processError(server, pluginConfig, error) - server.statsd?.increment(`plugin.retry.${retry.type}.${retry.pluginConfigId}.ERROR`) - } - } else { - server.statsd?.increment(`plugin.retry.${retry.type}.${retry.pluginConfigId}.SKIP`) - Sentry.captureMessage(`Retrying for plugin config ${retry.pluginConfigId} that does not exist`, { - extra: { retry: JSON.stringify(retry) }, - }) - } - server.statsd?.timing(`plugin.retry.${retry.type}.${retry.pluginConfigId}`, timer) - return response -} diff --git a/src/worker/plugins/setup.ts b/src/worker/plugins/setup.ts index ca147d3a..6918b727 100644 --- a/src/worker/plugins/setup.ts +++ b/src/worker/plugins/setup.ts @@ -1,6 +1,6 @@ import { PluginAttachment } from '@posthog/plugin-scaffold' -import { Plugin, PluginConfig, PluginConfigId, PluginId, PluginsServer, TeamId } from '../../types' +import { Plugin, PluginConfig, PluginConfigId, PluginId, PluginsServer, PluginTaskType, TeamId } from '../../types' import { getPluginAttachmentRows, getPluginConfigRows, getPluginRows } from '../../utils/db/sql' import { status } from '../../utils/status' import { LazyPluginVM } from '../vm/lazy' @@ -113,7 +113,7 @@ export async function loadSchedule(server: PluginsServer): Promise { let count = 0 for (const [id, pluginConfig] of server.pluginConfigs) { - const tasks = (await pluginConfig.vm?.getTasks()) ?? {} + const tasks = (await pluginConfig.vm?.getTasks(PluginTaskType.Schedule)) ?? {} for (const [taskName, task] of Object.entries(tasks)) { if (task && taskName in pluginSchedule) { pluginSchedule[taskName].push(id) diff --git a/src/worker/tasks.ts b/src/worker/tasks.ts new file mode 100644 index 00000000..b27aa939 --- /dev/null +++ b/src/worker/tasks.ts @@ -0,0 +1,51 @@ +import { PluginEvent } from '@posthog/plugin-scaffold/src/types' + +import { EnqueuedJob, PluginsServer, PluginTaskType } from '../types' +import { ingestEvent } from './ingestion/ingest-event' +import { runPlugins, runPluginsOnBatch, runPluginTask } from './plugins/run' +import { loadSchedule, setupPlugins } from './plugins/setup' +import { teardownPlugins } from './plugins/teardown' + +type TaskRunner = (server: PluginsServer, args: any) => Promise | any + +export const workerTasks: Record = { + hello: (server, args) => { + return `hello ${args}!` + }, + processEvent: (server, args: { event: PluginEvent }) => { + return runPlugins(server, args.event) + }, + processEventBatch: (server, args: { batch: PluginEvent[] }) => { + return runPluginsOnBatch(server, args.batch) + }, + runJob: (server, { job }: { job: EnqueuedJob }) => { + return runPluginTask(server, job.type, PluginTaskType.Job, job.pluginConfigId, job.payload) + }, + runEveryMinute: (server, args: { pluginConfigId: number }) => { + return runPluginTask(server, 'runEveryMinute', PluginTaskType.Schedule, args.pluginConfigId) + }, + runEveryHour: (server, args: { pluginConfigId: number }) => { + return runPluginTask(server, 'runEveryHour', PluginTaskType.Schedule, args.pluginConfigId) + }, + runEveryDay: (server, args: { pluginConfigId: number }) => { + return runPluginTask(server, 'runEveryDay', PluginTaskType.Schedule, args.pluginConfigId) + }, + getPluginSchedule: (server) => { + return server.pluginSchedule + }, + ingestEvent: async (server, args: { event: PluginEvent }) => { + return await ingestEvent(server, args.event) + }, + reloadPlugins: async (server) => { + await setupPlugins(server) + }, + reloadSchedule: async (server) => { + await loadSchedule(server) + }, + teardownPlugins: async (server) => { + await teardownPlugins(server) + }, + flushKafkaMessages: async (server) => { + await server.kafkaProducer?.flush() + }, +} diff --git a/src/worker/vm/extensions/jobs.ts b/src/worker/vm/extensions/jobs.ts new file mode 100644 index 00000000..74153ef4 --- /dev/null +++ b/src/worker/vm/extensions/jobs.ts @@ -0,0 +1,72 @@ +import { PluginConfig, PluginsServer } from '../../../types' + +type JobRunner = { + runAt: (date: Date) => Promise + runIn: (duration: number, unit: string) => Promise + runNow: () => Promise +} +type Job = (payload?: any) => JobRunner +type Jobs = Record + +const milliseconds = 1 +const seconds = 1000 * milliseconds +const minutes = 60 * seconds +const hours = 60 * minutes +const days = 24 * hours +const weeks = 7 * days +const months = 30 * weeks +const quarters = 13 * weeks +const years = 365 * days +const durations: Record = { + milliseconds, + seconds, + minutes, + hours, + days, + weeks, + months, + quarters, + years, +} + +export function durationToMs(duration: number, unit: string): number { + unit = `${unit}${unit.endsWith('s') ? '' : 's'}` + if (typeof durations[unit] === 'undefined') { + throw new Error(`Unknown time unit: ${unit}`) + } + return durations[unit] * duration +} + +export function createJobs(server: PluginsServer, pluginConfig: PluginConfig): Jobs { + const runJob = async (type: string, payload: Record, timestamp: number) => { + await server.jobQueueManager.enqueue({ + type, + payload, + timestamp, + pluginConfigId: pluginConfig.id, + pluginConfigTeam: pluginConfig.team_id, + }) + } + + return new Proxy( + {}, + { + get(target, key) { + return function createTaskRunner(payload: Record): JobRunner { + return { + runAt: async function runAt(date: Date) { + await runJob(key.toString(), payload, date.valueOf()) + }, + runIn: async function runIn(duration, unit) { + const timestamp = new Date().valueOf() + durationToMs(duration, unit) + await runJob(key.toString(), payload, timestamp) + }, + runNow: async function runNow() { + await runJob(key.toString(), payload, new Date().valueOf()) + }, + } + } + }, + } + ) +} diff --git a/src/worker/vm/extensions/retry.ts b/src/worker/vm/extensions/retry.ts deleted file mode 100644 index 9d4899f3..00000000 --- a/src/worker/vm/extensions/retry.ts +++ /dev/null @@ -1,23 +0,0 @@ -import { PluginConfig, PluginsServer } from '../../../types' - -const minRetry = process.env.NODE_ENV === 'test' ? 1 : 30 - -// TODO: add type to scaffold -export function createRetry( - server: PluginsServer, - pluginConfig: PluginConfig -): (type: string, payload: any, retry_in?: number) => Promise { - return async (type: string, payload: any, retry_in = 30) => { - if (retry_in < minRetry || retry_in > 86400) { - throw new Error(`Retries must happen between ${minRetry} seconds and 24 hours from now`) - } - const timestamp = new Date().valueOf() + retry_in * 1000 - await server.retryQueueManager.enqueue({ - type, - payload, - timestamp, - pluginConfigId: pluginConfig.id, - pluginConfigTeam: pluginConfig.team_id, - }) - } -} diff --git a/src/worker/vm/lazy.ts b/src/worker/vm/lazy.ts index 7829798d..a64168e9 100644 --- a/src/worker/vm/lazy.ts +++ b/src/worker/vm/lazy.ts @@ -5,6 +5,7 @@ import { PluginLogEntryType, PluginsServer, PluginTask, + PluginTaskType, } from '../../types' import { clearError, processError } from '../../utils/db/error' import { status } from '../../utils/status' @@ -70,15 +71,11 @@ export class LazyPluginVM { return (await this.resolveInternalVm)?.methods.teardownPlugin || null } - async getOnRetry(): Promise { - return (await this.resolveInternalVm)?.methods.onRetry || null + async getTask(name: string, type: PluginTaskType): Promise { + return (await this.resolveInternalVm)?.tasks?.[type]?.[name] || null } - async getTask(name: string): Promise { - return (await this.resolveInternalVm)?.tasks[name] || null - } - - async getTasks(): Promise> { - return (await this.resolveInternalVm)?.tasks || {} + async getTasks(type: PluginTaskType): Promise> { + return (await this.resolveInternalVm)?.tasks?.[type] || {} } } diff --git a/src/worker/vm/vm.ts b/src/worker/vm/vm.ts index 41d3bad6..4ce7659d 100644 --- a/src/worker/vm/vm.ts +++ b/src/worker/vm/vm.ts @@ -6,8 +6,8 @@ import { createCache } from './extensions/cache' import { createConsole } from './extensions/console' import { createGeoIp } from './extensions/geoip' import { createGoogle } from './extensions/google' +import { createJobs } from './extensions/jobs' import { createPosthog } from './extensions/posthog' -import { createRetry } from './extensions/retry' import { createStorage } from './extensions/storage' import { imports } from './imports' import { transformCode } from './transforms' @@ -66,7 +66,7 @@ export async function createPluginConfigVM( attachments: pluginConfig.attachments, storage: createStorage(server, pluginConfig), geoip: createGeoIp(server), - retry: createRetry(server, pluginConfig), + jobs: createJobs(server, pluginConfig), }, '__pluginHostMeta' ) @@ -154,17 +154,31 @@ export async function createPluginConfigVM( teardownPlugin: __asyncFunctionGuard(__bindMeta('teardownPlugin')), processEvent: __asyncFunctionGuard(__bindMeta('processEvent')), processEventBatch: __asyncFunctionGuard(__bindMeta('processEventBatch')), - onRetry: __asyncFunctionGuard(__bindMeta('onRetry')), }; - // gather the runEveryX commands and export in __tasks - const __tasks = {}; + const __tasks = { + schedule: {}, + job: {}, + }; + for (const exportDestination of __getExportDestinations().reverse()) { + // gather the runEveryX commands and export in __tasks for (const [name, value] of Object.entries(exportDestination)) { if (name.startsWith("runEvery") && typeof value === 'function') { - __tasks[name] = { + __tasks.schedule[name] = { name: name, - type: 'runEvery', + type: 'schedule', + exec: __bindMeta(value) + } + } + } + + // gather all jobs + if (typeof exportDestination['jobs'] === 'object') { + for (const [key, value] of Object.entries(exportDestination['jobs'])) { + __tasks.job[key] = { + name: key, + type: 'job', exec: __bindMeta(value) } } diff --git a/src/worker/worker.ts b/src/worker/worker.ts index 7a198285..4740776d 100644 --- a/src/worker/worker.ts +++ b/src/worker/worker.ts @@ -6,14 +6,12 @@ import { processError } from '../utils/db/error' import { createServer } from '../utils/db/server' import { status } from '../utils/status' import { cloneObject, pluginConfigIdFromStack } from '../utils/utils' -import { ingestEvent } from './ingestion/ingest-event' -import { runOnRetry, runPlugins, runPluginsOnBatch, runPluginTask } from './plugins/run' -import { loadSchedule, setupPlugins } from './plugins/setup' -import { teardownPlugins } from './plugins/teardown' +import { setupPlugins } from './plugins/setup' +import { workerTasks } from './tasks' -type TaskWorker = ({ task, args }: { task: string; args: any }) => Promise +export type PiscinaTaskWorker = ({ task, args }: { task: string; args: any }) => Promise -export async function createWorker(config: PluginsServerConfig, threadId: number): Promise { +export async function createWorker(config: PluginsServerConfig, threadId: number): Promise { initApp(config) status.info('๐Ÿงต', `Starting Piscina worker thread ${threadId}โ€ฆ`) @@ -30,49 +28,23 @@ export async function createWorker(config: PluginsServerConfig, threadId: number return createTaskRunner(server) } -export const createTaskRunner = (server: PluginsServer): TaskWorker => async ({ task, args }) => { +export const createTaskRunner = (server: PluginsServer): PiscinaTaskWorker => async ({ task, args }) => { const timer = new Date() let response Sentry.setContext('task', { task, args }) - if (task === 'hello') { - response = `hello ${args[0]}!` - } - if (task === 'processEvent') { - const processedEvent = await runPlugins(server, args.event) - // must clone the object, as we may get from VM2 something like { ..., properties: Proxy {} } - response = cloneObject(processedEvent as Record) - } - if (task === 'processEventBatch') { - const processedEvents = await runPluginsOnBatch(server, args.batch) - // must clone the object, as we may get from VM2 something like { ..., properties: Proxy {} } - response = cloneObject(processedEvents as any[]) - } - if (task === 'retry') { - response = await runOnRetry(server, args.retry) - } - if (task === 'getPluginSchedule') { - response = cloneObject(server.pluginSchedule) - } - if (task === 'ingestEvent') { - response = cloneObject(await ingestEvent(server, args.event)) - } - if (task.startsWith('runEvery')) { - const { pluginConfigId } = args - response = cloneObject(await runPluginTask(server, task, pluginConfigId)) - } - if (task === 'reloadPlugins') { - await setupPlugins(server) - } - if (task === 'reloadSchedule') { - await loadSchedule(server) - } - if (task === 'teardownPlugins') { - await teardownPlugins(server) - } - if (task === 'flushKafkaMessages') { - await server.kafkaProducer?.flush() + if (task in workerTasks) { + try { + // must clone the object, as we may get from VM2 something like { ..., properties: Proxy {} } + response = cloneObject(await workerTasks[task](server, args)) + } catch (e) { + status.info('๐Ÿ””', e) + Sentry.captureException(e) + response = { error: e.message } + } + } else { + response = { error: `Worker task "${task}" not found in: ${Object.keys(workerTasks).join(', ')}` } } server.statsd?.timing(`piscina_task.${task}`, timer) diff --git a/tests/errors.test.ts b/tests/errors.test.ts index fc003576..e69cd3f3 100644 --- a/tests/errors.test.ts +++ b/tests/errors.test.ts @@ -23,11 +23,13 @@ describe('error do not take down ingestion', () => { beforeEach(async () => { await resetTestDatabase(` - async function processEvent (event, meta) { - if (event.properties.crash === 'await') { - await meta.retry('test', {}, -100) - } else if (event.properties.crash === 'void') { - void meta.retry('test', {}, -100) + export async function processEvent (event, { jobs }) { + if (event.properties.crash === 'throw') { + throw new Error('error thrown in plugin') + } else if (event.properties.crash === 'throw in promise') { + void new Promise(() => { throw new Error('error thrown in plugin') }).then(() => {}) + } else if (event.properties.crash === 'reject in promise') { + void new Promise((_, rejects) => { rejects(new Error('error thrown in plugin')) }).then(() => {}) } return event } @@ -47,12 +49,12 @@ describe('error do not take down ingestion', () => { await stopServer() }) - test('awaited errors', async () => { + test('thrown errors', async () => { expect((await server.db.fetchEvents()).length).toBe(0) expect(await getErrorForPluginConfig(pluginConfig39.id)).toBe(null) for (let i = 0; i < 4; i++) { - posthog.capture('broken event', { crash: 'await' }) + posthog.capture('broken event', { crash: 'throw' }) } await delayUntilEventIngested(() => server.db.fetchEvents(), 4) @@ -60,7 +62,23 @@ describe('error do not take down ingestion', () => { expect((await server.db.fetchEvents()).length).toBe(4) const error2 = await getErrorForPluginConfig(pluginConfig39.id) - expect(error2.message).toBe('Retries must happen between 1 seconds and 24 hours from now') + expect(error2.message).toBe('error thrown in plugin') + }) + + test('unhandled promise errors', async () => { + expect((await server.db.fetchEvents()).length).toBe(0) + expect(await getErrorForPluginConfig(pluginConfig39.id)).toBe(null) + + for (let i = 0; i < 4; i++) { + posthog.capture('broken event', { crash: 'throw in promise' }) + } + + await delayUntilEventIngested(() => server.db.fetchEvents(), 4) + + expect((await server.db.fetchEvents()).length).toBe(4) + + const error2 = await getErrorForPluginConfig(pluginConfig39.id) + expect(error2.message).toBe('error thrown in plugin') }) test('unhandled promise rejections', async () => { @@ -68,7 +86,7 @@ describe('error do not take down ingestion', () => { expect(await getErrorForPluginConfig(pluginConfig39.id)).toBe(null) for (let i = 0; i < 4; i++) { - posthog.capture('broken event', { crash: 'void' }) + posthog.capture('broken event', { crash: 'reject in promise' }) } await delayUntilEventIngested(() => server.db.fetchEvents(), 4) @@ -76,6 +94,6 @@ describe('error do not take down ingestion', () => { expect((await server.db.fetchEvents()).length).toBe(4) const error2 = await getErrorForPluginConfig(pluginConfig39.id) - expect(error2.message).toBe('Retries must happen between 1 seconds and 24 hours from now') + expect(error2.message).toBe('error thrown in plugin') }) }) diff --git a/tests/helpers/graphile.ts b/tests/helpers/graphile.ts index 6cd82ed5..f2033ea7 100644 --- a/tests/helpers/graphile.ts +++ b/tests/helpers/graphile.ts @@ -1,22 +1,24 @@ import { makeWorkerUtils } from 'graphile-worker' import { Pool } from 'pg' -import { defaultConfig } from '../../src/config/config' -import { status } from '../../src/utils/status' +import { PluginsServerConfig } from '../../src/types' -export async function resetGraphileSchema(): Promise { - const db = new Pool({ connectionString: defaultConfig.DATABASE_URL }) +export async function resetGraphileSchema(server: PluginsServerConfig): Promise { + const graphileUrl = server.JOB_QUEUE_GRAPHILE_URL || server.DATABASE_URL + const db = new Pool({ connectionString: graphileUrl }) try { await db.query('DROP SCHEMA graphile_worker CASCADE') - } catch (e) { - status.error('๐Ÿ˜ฑ', `Could not dump graphile_worker schema: ${e.message}`) + } catch (error) { + if (error.message !== 'schema "graphile_worker" does not exist') { + throw error + } } finally { await db.end() } const workerUtils = await makeWorkerUtils({ - connectionString: defaultConfig.DATABASE_URL, + connectionString: graphileUrl, }) await workerUtils.migrate() await workerUtils.release() diff --git a/tests/jobs.test.ts b/tests/jobs.test.ts index 2e044171..a13bdefb 100644 --- a/tests/jobs.test.ts +++ b/tests/jobs.test.ts @@ -1,10 +1,11 @@ +import { defaultConfig } from '../src/config/config' import { LOCKED_RESOURCE } from '../src/main/job-queues/job-queue-consumer' -import { startPluginsServer } from '../src/main/pluginsServer' -import { LogLevel } from '../src/types' +import { ServerInstance, startPluginsServer } from '../src/main/pluginsServer' +import { LogLevel, PluginsServerConfig } from '../src/types' import { createServer } from '../src/utils/db/server' import { delay } from '../src/utils/utils' import { makePiscina } from '../src/worker/piscina' -import { createPosthog } from '../src/worker/vm/extensions/posthog' +import { createPosthog, DummyPostHog } from '../src/worker/vm/extensions/posthog' import { imports } from '../src/worker/vm/imports' import { resetGraphileSchema } from './helpers/graphile' import { pluginConfig39 } from './helpers/plugins' @@ -15,90 +16,105 @@ jest.setTimeout(60000) // 60 sec timeout const { console: testConsole } = imports['test-utils/write-to-file'] +const testCode = ` + import { console } from 'test-utils/write-to-file' + + export const jobs = { + logReply: (text, meta) => { + console.log('reply', text) + } + } + export async function processEvent (event, { jobs }) { + console.log('processEvent') + if (event.properties?.type === 'runIn') { + jobs.logReply('runIn').runIn(1, 'second') + } else if (event.properties?.type === 'runAt') { + jobs.logReply('runAt').runAt(new Date()) + } else if (event.properties?.type === 'runNow') { + jobs.logReply('runNow').runNow() + } + return event + } +` + +const createConfig = (jobQueues: string): PluginsServerConfig => ({ + ...defaultConfig, + WORKER_CONCURRENCY: 2, + LOG_LEVEL: LogLevel.Debug, + JOB_QUEUES: jobQueues, +}) + +async function waitForLogEntries(number: number) { + const timeout = 20000 + const start = new Date().valueOf() + while (testConsole.read().length < number) { + await delay(200) + if (new Date().valueOf() - start > timeout) { + console.error(`Did not find ${number} console logs:`, testConsole.read()) + throw new Error(`Did not get ${number} console logs within ${timeout / 1000} seconds`) + } + } +} + describe('job queues', () => { + let server: ServerInstance + let posthog: DummyPostHog + beforeEach(async () => { testConsole.reset() - const [server, stopServer] = await createServer() - const redis = await server.redisPool.acquire() + // reset lock in redis + const [tempServer, stopTempServer] = await createServer() + const redis = await tempServer.redisPool.acquire() await redis.del(LOCKED_RESOURCE) - await server.redisPool.release(redis) - await stopServer() + await tempServer.redisPool.release(redis) + await stopTempServer() + + // reset test code + await resetTestDatabase(testCode) + }) + + afterEach(async () => { + await server.stop() }) describe('fs queue', () => { - test('onRetry gets called', async () => { - const testCode = ` - import { console } from 'test-utils/write-to-file' - - export async function onRetry (type, payload, meta) { - console.log('retrying event!', type) - } - export async function processEvent (event, meta) { - if (event.properties?.hi === 'ha') { - console.log('processEvent') - meta.retry('processEvent', event, 1) - } - return event - } - ` - await resetTestDatabase(testCode) - const server = await startPluginsServer( - { - WORKER_CONCURRENCY: 2, - LOG_LEVEL: LogLevel.Debug, - RETRY_QUEUES: 'fs', - }, - makePiscina - ) - const posthog = createPosthog(server.server, pluginConfig39) - - posthog.capture('my event', { hi: 'ha' }) - await delay(10000) - - expect(testConsole.read()).toEqual([['processEvent'], ['retrying event!', 'processEvent']]) - - await server.stop() + beforeEach(async () => { + server = await startPluginsServer(createConfig('fs'), makePiscina) + posthog = createPosthog(server.server, pluginConfig39) + }) + + test('jobs get scheduled with runIn', async () => { + posthog.capture('my event', { type: 'runIn' }) + await waitForLogEntries(2) + expect(testConsole.read()).toEqual([['processEvent'], ['reply', 'runIn']]) + }) + + test('jobs get scheduled with runAt', async () => { + posthog.capture('my event', { type: 'runAt' }) + await waitForLogEntries(2) + expect(testConsole.read()).toEqual([['processEvent'], ['reply', 'runAt']]) + }) + + test('jobs get scheduled with runNow', async () => { + posthog.capture('my event', { type: 'runNow' }) + await waitForLogEntries(2) + expect(testConsole.read()).toEqual([['processEvent'], ['reply', 'runNow']]) }) }) describe('graphile', () => { beforeEach(async () => { - await resetGraphileSchema() + const config = createConfig('graphile') + await resetGraphileSchema(config) + server = await startPluginsServer(config, makePiscina) + posthog = createPosthog(server.server, pluginConfig39) }) - test('graphile retry queue', async () => { - const testCode = ` - import { console } from 'test-utils/write-to-file' - - export async function onRetry (type, payload, meta) { - console.log('retrying event!', type) - } - export async function processEvent (event, meta) { - if (event.properties?.hi === 'ha') { - console.log('processEvent') - meta.retry('processEvent', event, 1) - } - return event - } - ` - await resetTestDatabase(testCode) - const server = await startPluginsServer( - { - WORKER_CONCURRENCY: 2, - LOG_LEVEL: LogLevel.Debug, - RETRY_QUEUES: 'graphile', - }, - makePiscina - ) - const posthog = createPosthog(server.server, pluginConfig39) - - posthog.capture('my event', { hi: 'ha' }) - await delay(5000) - - expect(testConsole.read()).toEqual([['processEvent'], ['retrying event!', 'processEvent']]) - - await server.stop() + test('graphile job queue', async () => { + posthog.capture('my event', { type: 'runIn' }) + await waitForLogEntries(2) + expect(testConsole.read()).toEqual([['processEvent'], ['reply', 'runIn']]) }) }) }) diff --git a/tests/plugins.test.ts b/tests/plugins.test.ts index 553f0fa5..9c154b34 100644 --- a/tests/plugins.test.ts +++ b/tests/plugins.test.ts @@ -1,7 +1,7 @@ import { PluginEvent } from '@posthog/plugin-scaffold/src/types' import { mocked } from 'ts-jest/utils' -import { LogLevel, PluginsServer } from '../src/types' +import { LogLevel, PluginsServer, PluginTaskType } from '../src/types' import { clearError, processError } from '../src/utils/db/error' import { createServer } from '../src/utils/db/server' import { loadPlugin } from '../src/worker/plugins/loadPlugin' @@ -15,7 +15,7 @@ import { pluginAttachment1, pluginConfig39, } from './helpers/plugins' -import { getPluginAttachmentRows, getPluginConfigRows, getPluginRows, setError } from './helpers/sqlMock' +import { getPluginAttachmentRows, getPluginConfigRows, getPluginRows } from './helpers/sqlMock' jest.mock('../src/utils/db/sql') jest.mock('../src/utils/status') @@ -70,7 +70,6 @@ test('setupPlugins and runPlugins', async () => { expect(pluginConfig.vm).toBeDefined() const vm = await pluginConfig.vm!.resolveInternalVm expect(Object.keys(vm!.methods).sort()).toEqual([ - 'onRetry', 'processEvent', 'processEventBatch', 'setupPlugin', @@ -128,7 +127,7 @@ test('plugin meta has what it should have', async () => { 'config', 'geoip', 'global', - 'retry', + 'jobs', 'storage', ]) expect(returnedEvent!.properties!['attachments']).toEqual({ @@ -155,7 +154,7 @@ test('archive plugin with broken index.js does not do much', async () => { const { pluginConfigs } = mockServer const pluginConfig = pluginConfigs.get(39)! - expect(await pluginConfigs.get(39)!.vm!.getTasks()).toEqual({}) + expect(await pluginConfigs.get(39)!.vm!.getTasks(PluginTaskType.Schedule)).toEqual({}) const event = { event: '$test', properties: {}, team_id: 2 } as PluginEvent const returnedEvent = await runPlugins(mockServer, { ...event }) @@ -180,7 +179,7 @@ test('local plugin with broken index.js does not do much', async () => { const { pluginConfigs } = mockServer const pluginConfig = pluginConfigs.get(39)! - expect(await pluginConfigs.get(39)!.vm!.getTasks()).toEqual({}) + expect(await pluginConfigs.get(39)!.vm!.getTasks(PluginTaskType.Schedule)).toEqual({}) const event = { event: '$test', properties: {}, team_id: 2 } as PluginEvent const returnedEvent = await runPlugins(mockServer, { ...event }) @@ -211,7 +210,7 @@ test('plugin throwing error does not prevent ingestion and failure is noted in e await setupPlugins(mockServer) const { pluginConfigs } = mockServer - expect(await pluginConfigs.get(39)!.vm!.getTasks()).toEqual({}) + expect(await pluginConfigs.get(39)!.vm!.getTasks(PluginTaskType.Schedule)).toEqual({}) const event = { event: '$test', properties: {}, team_id: 2 } as PluginEvent const returnedEvent = await runPlugins(mockServer, { ...event }) @@ -244,7 +243,7 @@ test('events have property $plugins_succeeded set to the plugins that succeeded' await setupPlugins(mockServer) const { pluginConfigs } = mockServer - expect(await pluginConfigs.get(39)!.vm!.getTasks()).toEqual({}) + expect(await pluginConfigs.get(39)!.vm!.getTasks(PluginTaskType.Schedule)).toEqual({}) const event = { event: '$test', properties: {}, team_id: 2 } as PluginEvent const returnedEvent = await runPlugins(mockServer, { ...event }) @@ -282,7 +281,7 @@ test('archive plugin with broken plugin.json does not do much', async () => { `Can not load plugin.json for plugin test-maxmind-plugin ID ${plugin60.id} (organization ID ${commonOrganizationId})` ) - expect(await pluginConfigs.get(39)!.vm!.getTasks()).toEqual({}) + expect(await pluginConfigs.get(39)!.vm!.getTasks(PluginTaskType.Schedule)).toEqual({}) }) test('local plugin with broken plugin.json does not do much', async () => { @@ -306,7 +305,7 @@ test('local plugin with broken plugin.json does not do much', async () => { pluginConfigs.get(39)!, expect.stringContaining('Could not load posthog config at ') ) - expect(await pluginConfigs.get(39)!.vm!.getTasks()).toEqual({}) + expect(await pluginConfigs.get(39)!.vm!.getTasks(PluginTaskType.Schedule)).toEqual({}) unlink() }) @@ -329,7 +328,7 @@ test('plugin with http urls must have an archive', async () => { pluginConfigs.get(39)!, `Tried using undownloaded remote plugin test-maxmind-plugin ID ${plugin60.id} (organization ID ${commonOrganizationId} - global), which is not supported!` ) - expect(await pluginConfigs.get(39)!.vm!.getTasks()).toEqual({}) + expect(await pluginConfigs.get(39)!.vm!.getTasks(PluginTaskType.Schedule)).toEqual({}) }) test("plugin with broken archive doesn't load", async () => { @@ -350,7 +349,7 @@ test("plugin with broken archive doesn't load", async () => { pluginConfigs.get(39)!, Error('Could not read archive as .zip or .tgz') ) - expect(await pluginConfigs.get(39)!.vm!.getTasks()).toEqual({}) + expect(await pluginConfigs.get(39)!.vm!.getTasks(PluginTaskType.Schedule)).toEqual({}) }) test('plugin config order', async () => { diff --git a/tests/postgres/vm.lazy.test.ts b/tests/postgres/vm.lazy.test.ts index e39199cd..57a29280 100644 --- a/tests/postgres/vm.lazy.test.ts +++ b/tests/postgres/vm.lazy.test.ts @@ -1,5 +1,6 @@ import { mocked } from 'ts-jest/utils' +import { PluginTaskType } from '../../src/types' import { clearError, processError } from '../../src/utils/db/error' import { status } from '../../src/utils/status' import { LazyPluginVM } from '../../src/worker/vm/lazy' @@ -21,7 +22,9 @@ describe('LazyPluginVM', () => { processEvent: 'processEvent', }, tasks: { - runEveryMinute: 'runEveryMinute', + schedule: { + runEveryMinute: 'runEveryMinute', + }, }, } @@ -35,9 +38,9 @@ describe('LazyPluginVM', () => { expect(await vm.getProcessEvent()).toEqual('processEvent') expect(await vm.getProcessEventBatch()).toEqual(null) - expect(await vm.getTask('someTask')).toEqual(null) - expect(await vm.getTask('runEveryMinute')).toEqual('runEveryMinute') - expect(await vm.getTasks()).toEqual(mockVM.tasks) + expect(await vm.getTask('someTask', PluginTaskType.Schedule)).toEqual(null) + expect(await vm.getTask('runEveryMinute', PluginTaskType.Schedule)).toEqual('runEveryMinute') + expect(await vm.getTasks(PluginTaskType.Schedule)).toEqual(mockVM.tasks.schedule) }) it('logs info and clears errors on success', async () => { @@ -63,8 +66,8 @@ describe('LazyPluginVM', () => { expect(await vm.getProcessEvent()).toEqual(null) expect(await vm.getProcessEventBatch()).toEqual(null) - expect(await vm.getTask('runEveryMinute')).toEqual(null) - expect(await vm.getTasks()).toEqual({}) + expect(await vm.getTask('runEveryMinute', PluginTaskType.Schedule)).toEqual(null) + expect(await vm.getTasks(PluginTaskType.Schedule)).toEqual({}) }) it('logs failure', async () => { diff --git a/tests/postgres/vm.test.ts b/tests/postgres/vm.test.ts index 26b07684..7929b865 100644 --- a/tests/postgres/vm.test.ts +++ b/tests/postgres/vm.test.ts @@ -39,7 +39,6 @@ test('empty plugins', async () => { expect(Object.keys(vm).sort()).toEqual(['methods', 'tasks', 'vm']) expect(Object.keys(vm.methods).sort()).toEqual([ - 'onRetry', 'processEvent', 'processEventBatch', 'setupPlugin', @@ -725,10 +724,15 @@ test('runEvery', async () => { await resetTestDatabase(indexJs) const vm = await createPluginConfigVM(mockServer, pluginConfig39, indexJs) - expect(Object.keys(vm.tasks)).toEqual(['runEveryMinute', 'runEveryHour', 'runEveryDay']) - expect(Object.values(vm.tasks).map((v) => v?.name)).toEqual(['runEveryMinute', 'runEveryHour', 'runEveryDay']) - expect(Object.values(vm.tasks).map((v) => v?.type)).toEqual(['runEvery', 'runEvery', 'runEvery']) - expect(Object.values(vm.tasks).map((v) => typeof v?.exec)).toEqual(['function', 'function', 'function']) + expect(Object.keys(vm.tasks).sort()).toEqual(['job', 'schedule']) + expect(Object.keys(vm.tasks.schedule)).toEqual(['runEveryMinute', 'runEveryHour', 'runEveryDay']) + expect(Object.values(vm.tasks.schedule).map((v) => v?.name)).toEqual([ + 'runEveryMinute', + 'runEveryHour', + 'runEveryDay', + ]) + expect(Object.values(vm.tasks.schedule).map((v) => v?.type)).toEqual(['schedule', 'schedule', 'schedule']) + expect(Object.values(vm.tasks.schedule).map((v) => typeof v?.exec)).toEqual(['function', 'function', 'function']) }) test('runEvery must be a function', async () => { @@ -742,10 +746,10 @@ test('runEvery must be a function', async () => { await resetTestDatabase(indexJs) const vm = await createPluginConfigVM(mockServer, pluginConfig39, indexJs) - expect(Object.keys(vm.tasks)).toEqual(['runEveryMinute']) - expect(Object.values(vm.tasks).map((v) => v?.name)).toEqual(['runEveryMinute']) - expect(Object.values(vm.tasks).map((v) => v?.type)).toEqual(['runEvery']) - expect(Object.values(vm.tasks).map((v) => typeof v?.exec)).toEqual(['function']) + expect(Object.keys(vm.tasks.schedule)).toEqual(['runEveryMinute']) + expect(Object.values(vm.tasks.schedule).map((v) => v?.name)).toEqual(['runEveryMinute']) + expect(Object.values(vm.tasks.schedule).map((v) => v?.type)).toEqual(['schedule']) + expect(Object.values(vm.tasks.schedule).map((v) => typeof v?.exec)).toEqual(['function']) }) test('posthog in runEvery', async () => { @@ -760,7 +764,7 @@ test('posthog in runEvery', async () => { expect(Client).not.toHaveBeenCalled - const response = await vm.tasks.runEveryMinute.exec() + const response = await vm.tasks.schedule.runEveryMinute.exec() expect(response).toBe('haha') expect(Client).toHaveBeenCalledTimes(2) @@ -803,7 +807,7 @@ test('posthog in runEvery with timestamp', async () => { expect(Client).not.toHaveBeenCalled - const response = await vm.tasks.runEveryMinute.exec() + const response = await vm.tasks.schedule.runEveryMinute.exec() expect(response).toBe('haha') expect(Client).toHaveBeenCalledTimes(2) @@ -843,7 +847,7 @@ test('posthog.capture accepts user-defined distinct id', async () => { expect(Client).not.toHaveBeenCalled - const response = await vm.tasks.runEveryMinute.exec() + const response = await vm.tasks.schedule.runEveryMinute.exec() expect(response).toBe('haha') const mockClientInstance = (Client as any).mock.instances[1] diff --git a/tests/postgres/worker.test.ts b/tests/postgres/worker.test.ts index fc8b69c5..4ebc322b 100644 --- a/tests/postgres/worker.test.ts +++ b/tests/postgres/worker.test.ts @@ -281,7 +281,9 @@ describe('createTaskRunner()', () => { }) it('handles `runEvery` tasks', async () => { - mocked(runPluginTask).mockImplementation((server, task, pluginId) => Promise.resolve(`${task} for ${pluginId}`)) + mocked(runPluginTask).mockImplementation((server, task, taskType, pluginId) => + Promise.resolve(`${task} for ${pluginId}`) + ) expect(await taskRunner({ task: 'runEveryMinute', args: { pluginConfigId: 1 } })).toEqual( 'runEveryMinute for 1' diff --git a/tests/retry.test.ts b/tests/retry.test.ts deleted file mode 100644 index ae9f660c..00000000 --- a/tests/retry.test.ts +++ /dev/null @@ -1,104 +0,0 @@ -import { startPluginsServer } from '../src/main/pluginsServer' -import { LOCKED_RESOURCE } from '../src/main/services/retry-queue-consumer' -import { LogLevel } from '../src/types' -import { createServer } from '../src/utils/db/server' -import { delay } from '../src/utils/utils' -import { makePiscina } from '../src/worker/piscina' -import { createPosthog } from '../src/worker/vm/extensions/posthog' -import { imports } from '../src/worker/vm/imports' -import { resetGraphileSchema } from './helpers/graphile' -import { pluginConfig39 } from './helpers/plugins' -import { resetTestDatabase } from './helpers/sql' - -jest.mock('../src/utils/db/sql') -jest.setTimeout(60000) // 60 sec timeout - -const { console: testConsole } = imports['test-utils/write-to-file'] - -describe('retry queues', () => { - beforeEach(async () => { - testConsole.reset() - - const [server, stopServer] = await createServer() - const redis = await server.redisPool.acquire() - await redis.del(LOCKED_RESOURCE) - await server.redisPool.release(redis) - await stopServer() - }) - - describe('fs queue', () => { - test('onRetry gets called', async () => { - const testCode = ` - import { console } from 'test-utils/write-to-file' - - export async function onRetry (type, payload, meta) { - console.log('retrying event!', type) - } - export async function processEvent (event, meta) { - if (event.properties?.hi === 'ha') { - console.log('processEvent') - meta.retry('processEvent', event, 1) - } - return event - } - ` - await resetTestDatabase(testCode) - const server = await startPluginsServer( - { - WORKER_CONCURRENCY: 2, - LOG_LEVEL: LogLevel.Debug, - RETRY_QUEUES: 'fs', - }, - makePiscina - ) - const posthog = createPosthog(server.server, pluginConfig39) - - posthog.capture('my event', { hi: 'ha' }) - await delay(10000) - - expect(testConsole.read()).toEqual([['processEvent'], ['retrying event!', 'processEvent']]) - - await server.stop() - }) - }) - - describe('graphile', () => { - beforeEach(async () => { - await resetGraphileSchema() - }) - - test('graphile retry queue', async () => { - const testCode = ` - import { console } from 'test-utils/write-to-file' - - export async function onRetry (type, payload, meta) { - console.log('retrying event!', type) - } - export async function processEvent (event, meta) { - if (event.properties?.hi === 'ha') { - console.log('processEvent') - meta.retry('processEvent', event, 1) - } - return event - } - ` - await resetTestDatabase(testCode) - const server = await startPluginsServer( - { - WORKER_CONCURRENCY: 2, - LOG_LEVEL: LogLevel.Debug, - RETRY_QUEUES: 'graphile', - }, - makePiscina - ) - const posthog = createPosthog(server.server, pluginConfig39) - - posthog.capture('my event', { hi: 'ha' }) - await delay(5000) - - expect(testConsole.read()).toEqual([['processEvent'], ['retrying event!', 'processEvent']]) - - await server.stop() - }) - }) -}) diff --git a/tests/schedule.test.ts b/tests/schedule.test.ts index b9d78abc..10bfeb4a 100644 --- a/tests/schedule.test.ts +++ b/tests/schedule.test.ts @@ -3,7 +3,7 @@ import { PluginEvent } from '@posthog/plugin-scaffold/src/types' import { loadPluginSchedule, LOCKED_RESOURCE, - runTasksDebounced, + runScheduleDebounced, startSchedule, waitForTasksToFinish, } from '../src/main/services/schedule' @@ -30,7 +30,7 @@ function createEvent(index = 0): PluginEvent { } } -test('runTasksDebounced', async () => { +test('runScheduleDebounced', async () => { const workerThreads = 1 const testCode = ` const counterKey = 'test_counter_2' @@ -58,9 +58,9 @@ test('runTasksDebounced', async () => { const event1 = await processEvent(createEvent()) expect(event1.properties['counter']).toBe(0) - runTasksDebounced(server, piscina, 'runEveryMinute') - runTasksDebounced(server, piscina, 'runEveryMinute') - runTasksDebounced(server, piscina, 'runEveryMinute') + runScheduleDebounced(server, piscina, 'runEveryMinute') + runScheduleDebounced(server, piscina, 'runEveryMinute') + runScheduleDebounced(server, piscina, 'runEveryMinute') await delay(100) const event2 = await processEvent(createEvent()) @@ -77,7 +77,7 @@ test('runTasksDebounced', async () => { await closeServer() }) -test('runTasksDebounced exception', async () => { +test('runScheduleDebounced exception', async () => { const workerThreads = 2 const testCode = ` async function runEveryMinute (meta) { @@ -90,7 +90,7 @@ test('runTasksDebounced exception', async () => { const [server, closeServer] = await createServer({ LOG_LEVEL: LogLevel.Log }) server.pluginSchedule = await loadPluginSchedule(piscina) - runTasksDebounced(server, piscina, 'runEveryMinute') + runScheduleDebounced(server, piscina, 'runEveryMinute') await waitForTasksToFinish(server)