From e85dddd9bdbac3aa534d3f4812df12dcee69f1f7 Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Tue, 20 Apr 2021 11:54:26 +0200 Subject: [PATCH 01/22] extract redlock from schedule --- src/main/services/redlock.ts | 92 +++++++++++++++++++++++++++++++++++ src/main/services/schedule.ts | 86 ++++++-------------------------- 2 files changed, 107 insertions(+), 71 deletions(-) create mode 100644 src/main/services/redlock.ts diff --git a/src/main/services/redlock.ts b/src/main/services/redlock.ts new file mode 100644 index 00000000..fcf40df1 --- /dev/null +++ b/src/main/services/redlock.ts @@ -0,0 +1,92 @@ +import * as Sentry from '@sentry/node' +import Redlock from 'redlock' + +import { status } from '../../shared/status' +import { createRedis } from '../../shared/utils' +import { PluginsServer } from '../../types' + +export async function startRedlock( + server: PluginsServer, + resource: string, + onLock: () => Promise | void, + onUnlock: () => Promise | void, + ttl = 60 +): 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 + 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 retires 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/schedule.ts b/src/main/services/schedule.ts index a1ac9ce3..6514e02f 100644 --- a/src/main/services/schedule.ts +++ b/src/main/services/schedule.ts @@ -1,12 +1,11 @@ import Piscina from '@posthog/piscina' -import * as Sentry from '@sentry/node' import * as schedule from 'node-schedule' -import Redlock from 'redlock' import { processError } from '../../shared/error' import { status } from '../../shared/status' -import { createRedis, delay } from '../../shared/utils' +import { delay } from '../../shared/utils' import { PluginConfigId, PluginsServer, ScheduleControl } from '../../types' +import { startRedlock } from './redlock' export const LOCKED_RESOURCE = 'plugin-server:locks:schedule' @@ -19,74 +18,9 @@ export async function startSchedule( let stopped = false let weHaveTheLock = false - let lock: Redlock.Lock - let lockTimeout: NodeJS.Timeout - - const lockTTL = server.SCHEDULE_LOCK_TTL * 1000 // 60 sec - 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 retires 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 client error occurred:\n', error) - Sentry.captureException(error) - }) - - const tryToGetTheLock = async () => { - try { - lock = await redlock.lock(LOCKED_RESOURCE, lockTTL) - weHaveTheLock = true - - status.info('🔒', 'Scheduler lock 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:\n', error) - Sentry.captureException(error) - weHaveTheLock = false - lockTimeout = setTimeout(tryToGetTheLock, 0) - } - } - - lockTimeout = setTimeout(extendLock, extendDelay) - - onLock?.() - } catch (error) { - if (stopped) { - return - } - weHaveTheLock = false - if (error instanceof Redlock.LockError) { - lockTimeout = setTimeout(tryToGetTheLock, retryDelay) - } else { - Sentry.captureException(error) - status.error('🔴', 'Redlock error:\n', error) - } - } - } - let pluginSchedulePromise = loadPluginSchedule(piscina) server.pluginSchedule = await pluginSchedulePromise - lockTimeout = setTimeout(tryToGetTheLock, 0) - const runEveryMinuteJob = schedule.scheduleJob('* * * * *', async () => { !stopped && weHaveTheLock && @@ -106,15 +40,25 @@ export async function startSchedule( runTasksDebounced(server!, piscina!, 'runEveryDay') }) + const unlock = await startRedlock( + server, + LOCKED_RESOURCE, + () => { + weHaveTheLock = true + }, + () => { + weHaveTheLock = false + }, + server.SCHEDULE_LOCK_TTL + ) + const stopSchedule = async () => { stopped = true - lockTimeout && clearTimeout(lockTimeout) runEveryDayJob && schedule.cancelJob(runEveryDayJob) runEveryHourJob && schedule.cancelJob(runEveryHourJob) runEveryMinuteJob && schedule.cancelJob(runEveryMinuteJob) - await lock?.unlock().catch(Sentry.captureException) - await redis.quit() + await unlock() await waitForTasksToFinish(server!) } From 7333c0ce107ca66d5ae6d057e5228a241731269f Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Tue, 20 Apr 2021 16:26:18 +0200 Subject: [PATCH 02/22] implement generic retrying --- src/main/pluginsServer.ts | 10 ++- src/main/queue.ts | 14 +++-- src/main/retry/fs-queue.ts | 74 +++++++++++++++++++++++ src/main/retry/retry-queue-manager.ts | 68 +++++++++++++++++++++ src/main/services/retry-queue-consumer.ts | 38 ++++++++++++ src/shared/config.ts | 2 + src/shared/server.ts | 2 + src/types.ts | 36 ++++++++--- src/worker/plugins/run.ts | 17 +++++- src/worker/vm/extensions/retry.ts | 23 +++++++ src/worker/vm/lazy.ts | 4 ++ src/worker/vm/vm.ts | 3 + src/worker/worker.ts | 5 +- tests/retry.test.ts | 53 ++++++++++++++++ 14 files changed, 332 insertions(+), 17 deletions(-) create mode 100644 src/main/retry/fs-queue.ts create mode 100644 src/main/retry/retry-queue-manager.ts create mode 100644 src/main/services/retry-queue-consumer.ts create mode 100644 src/worker/vm/extensions/retry.ts create mode 100644 tests/retry.test.ts diff --git a/src/main/pluginsServer.ts b/src/main/pluginsServer.ts index f817e9b1..cdb198d8 100644 --- a/src/main/pluginsServer.ts +++ b/src/main/pluginsServer.ts @@ -10,9 +10,10 @@ import { defaultConfig } from '../shared/config' import { createServer } from '../shared/server' import { status } from '../shared/status' import { createRedis, delay } from '../shared/utils' -import { PluginsServer, PluginsServerConfig, Queue, ScheduleControl } from '../types' +import { PluginsServer, PluginsServerConfig, Queue, RetryQueueConsumerControl, ScheduleControl } from '../types' import { createMmdbServer, performMmdbStalenessCheck, prepareMmdb } from './mmdb' import { startQueue } from './queue' +import { startRetryQueueConsumer } from './services/retry-queue-consumer' import { startSchedule } from './services/schedule' import { startFastifyInstance, stopFastifyInstance } from './web/server' @@ -46,6 +47,7 @@ export async function startPluginsServer( let statsJob: schedule.Job | undefined let piscina: Piscina | undefined let queue: Queue | undefined + let retryQueueConsumer: RetryQueueConsumerControl | undefined let closeServer: () => Promise | undefined let scheduleControl: ScheduleControl | undefined let mmdbServer: net.Server | undefined @@ -71,6 +73,7 @@ export async function startPluginsServer( await pubSub?.quit() pingJob && schedule.cancelJob(pingJob) statsJob && schedule.cancelJob(statsJob) + await retryQueueConsumer?.stop() await scheduleControl?.stopSchedule() await new Promise((resolve, reject) => !mmdbServer @@ -127,9 +130,12 @@ export async function startPluginsServer( } scheduleControl = await startSchedule(server, piscina) + retryQueueConsumer = await startRetryQueueConsumer(server, piscina) + queue = await startQueue(server, piscina) piscina.on('drain', () => { - queue?.resume() + void queue?.resume() + void retryQueueConsumer?.resume() }) // use one extra connection for redis pubsub diff --git a/src/main/queue.ts b/src/main/queue.ts index e257d473..29a31720 100644 --- a/src/main/queue.ts +++ b/src/main/queue.ts @@ -14,9 +14,13 @@ export type WorkerMethods = { ingestEvent: (event: PluginEvent) => Promise } -function pauseQueueIfWorkerFull(queue: Queue | undefined, server: PluginsServer, piscina?: Piscina) { - if (queue && (piscina?.queueSize || 0) > (server.WORKER_CONCURRENCY || 4) * (server.WORKER_CONCURRENCY || 4)) { - void queue.pause() +export function pauseQueueIfWorkerFull( + pause: undefined | (() => void | Promise), + server: PluginsServer, + piscina?: Piscina +) { + if (pause && (piscina?.queueSize || 0) > (server.WORKER_CONCURRENCY || 4) * (server.WORKER_CONCURRENCY || 4)) { + void pause() } } @@ -75,10 +79,10 @@ function startQueueRedis(server: PluginsServer, piscina: Piscina | undefined, wo ...data, } as PluginEvent) try { - pauseQueueIfWorkerFull(celeryQueue, server, piscina) + pauseQueueIfWorkerFull(() => celeryQueue.pause(), server, piscina) const processedEvent = await workerMethods.processEvent(event) if (processedEvent) { - pauseQueueIfWorkerFull(celeryQueue, server, piscina) + pauseQueueIfWorkerFull(() => celeryQueue.pause(), server, piscina) await workerMethods.ingestEvent(processedEvent) } } catch (e) { diff --git a/src/main/retry/fs-queue.ts b/src/main/retry/fs-queue.ts new file mode 100644 index 00000000..c1de91e0 --- /dev/null +++ b/src/main/retry/fs-queue.ts @@ -0,0 +1,74 @@ +import { EnqueuedRetry, OnRetryCallback, RetryQueue } from '../../types' +import Timeout = NodeJS.Timeout +import * as fs from 'fs' + +export class FsQueue implements RetryQueue { + paused: boolean + started: boolean + interval: Timeout | null + filename: string + + constructor(filename = '/tmp/fs-queue-file.txt') { + if (process.env.NODE_ENV !== 'test') { + throw new Error('Can not use FsQueue outside tests') + } + this.paused = false + this.started = false + this.interval = null + this.filename = filename + // this.queued = [] + } + + enqueue(retry: EnqueuedRetry): Promise | void { + fs.appendFileSync(this.filename, `${JSON.stringify(retry)}\n`) + } + + startConsumer(onRetry: OnRetryCallback): void { + fs.writeFileSync(this.filename, '') + this.started = true + this.interval = setInterval(() => { + if (this.paused) { + return + } + const timestamp = new Date().valueOf() + const queue = fs + .readFileSync(this.filename) + .toString() + .split('\n') + .map((s) => { + try { + return JSON.parse(s) as EnqueuedRetry + } catch (e) { + return null + } + }) + .filter((a) => !!a) as EnqueuedRetry[] + + 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) + } + }, 1000) + } + + stopConsumer(): void { + this.started = false + this.interval && clearInterval(this.interval) + fs.unlinkSync(this.filename) + } + + pauseConsumer(): void { + this.paused = true + } + + isConsumerPaused(): boolean { + return this.paused + } + + resumeConsumer(): void { + this.paused = false + } +} diff --git a/src/main/retry/retry-queue-manager.ts b/src/main/retry/retry-queue-manager.ts new file mode 100644 index 00000000..ced7635a --- /dev/null +++ b/src/main/retry/retry-queue-manager.ts @@ -0,0 +1,68 @@ +import * as Sentry from '@sentry/node' + +import { EnqueuedRetry, OnRetryCallback, PluginsServer, RetryQueue } from '../../types' +import { FsQueue } from './fs-queue' + +const queues = { + fs: () => new FsQueue(), +} + +export class RetryQueueManager implements RetryQueue { + pluginsServer: PluginsServer + retryQueues: RetryQueue[] + + constructor(pluginsServer: PluginsServer) { + this.pluginsServer = pluginsServer + + this.retryQueues = pluginsServer.RETRY_QUEUES.split(',') + .map((q) => q.trim()) + .map( + (queue): RetryQueue => { + if (queues[queue as keyof typeof queues]) { + return queues[queue as keyof typeof queues]() + } else { + throw new Error(`Unknown retry queue "${queue}"`) + } + } + ) + } + + async enqueue(retry: EnqueuedRetry): Promise { + for (const retryQueue of this.retryQueues) { + try { + await retryQueue.enqueue(retry) + return + } catch (error) { + // if one fails, take the next queue + Sentry.captureException(error, { + extra: { + retry: JSON.stringify(retry), + queue: retryQueue.toString(), + queues: this.retryQueues.map((q) => q.toString()), + }, + }) + } + } + throw new Error('No RetryQueue available') + } + + async startConsumer(onRetry: OnRetryCallback): Promise { + await Promise.all(this.retryQueues.map((r) => r.startConsumer(onRetry))) + } + + async stopConsumer(): Promise { + await Promise.all(this.retryQueues.map((r) => r.stopConsumer())) + } + + async pauseConsumer(): Promise { + await Promise.all(this.retryQueues.map((r) => r.pauseConsumer())) + } + + isConsumerPaused(): boolean { + return !!this.retryQueues.find((r) => r.isConsumerPaused()) + } + + async resumeConsumer(): Promise { + await Promise.all(this.retryQueues.map((r) => r.resumeConsumer())) + } +} diff --git a/src/main/services/retry-queue-consumer.ts b/src/main/services/retry-queue-consumer.ts new file mode 100644 index 00000000..11c710b9 --- /dev/null +++ b/src/main/services/retry-queue-consumer.ts @@ -0,0 +1,38 @@ +import Piscina from '@posthog/piscina' + +import { status } from '../../shared/status' +import { OnRetryCallback, PluginsServer, RetryQueueConsumerControl } from '../../types' +import { pauseQueueIfWorkerFull } from '../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, + LOCKED_RESOURCE, + async () => { + status.info('🔄', 'Retry queue consumer lock aquired') + await server.retryQueueManager.startConsumer(onRetry) + }, + async () => { + status.info('🔄', 'Stopping retry queue consumer') + await server.retryQueueManager.stopConsumer() + }, + server.SCHEDULE_LOCK_TTL + ) + + return { stop: () => unlock(), resume: () => server.retryQueueManager.resumeConsumer() } +} diff --git a/src/shared/config.ts b/src/shared/config.ts index 1b149bc8..44aa03cb 100644 --- a/src/shared/config.ts +++ b/src/shared/config.ts @@ -58,6 +58,7 @@ export function getDefaultConfig(): PluginsServerConfig { DISTINCT_ID_LRU_SIZE: 10000, INTERNAL_MMDB_SERVER_PORT: 0, PLUGIN_SERVER_IDLE: false, + RETRY_QUEUES: '', } } @@ -99,6 +100,7 @@ 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: 'queue system and fallbacks to use for retries', } } diff --git a/src/shared/server.ts b/src/shared/server.ts index e622d910..18cab2db 100644 --- a/src/shared/server.ts +++ b/src/shared/server.ts @@ -10,6 +10,7 @@ import * as path from 'path' import { types as pgTypes } from 'pg' import { ConnectionOptions } from 'tls' +import { RetryQueueManager } from '../main/retry/retry-queue-manager' import { PluginsServer, PluginsServerConfig } from '../types' import { EventsProcessor } from '../worker/ingestion/process-event' import { defaultConfig } from './config' @@ -159,6 +160,7 @@ export async function createServer( // :TODO: This is only used on worker threads, not main server.eventsProcessor = new EventsProcessor(server as PluginsServer) + server.retryQueueManager = new RetryQueueManager(server as PluginsServer) const closeServer = async () => { server.mmdbUpdateJob?.cancel() diff --git a/src/types.ts b/src/types.ts index 0e6299f2..4b5bb5c1 100644 --- a/src/types.ts +++ b/src/types.ts @@ -1,4 +1,3 @@ -import { ReaderModel } from '@maxmind/geoip2-node' import ClickHouse from '@posthog/clickhouse' import { PluginAttachment, PluginConfigSchema, PluginEvent, Properties } from '@posthog/plugin-scaffold' import { Pool as GenericPool } from 'generic-pool' @@ -6,7 +5,6 @@ import { StatsD } from 'hot-shots' import { Redis } from 'ioredis' import { Kafka } from 'kafkajs' import { DateTime } from 'luxon' -import { Job } from 'node-schedule' import { Pool } from 'pg' import { VM } from 'vm2' @@ -72,6 +70,7 @@ export interface PluginsServerConfig extends Record { DISTINCT_ID_LRU_SIZE: number INTERNAL_MMDB_SERVER_PORT: number PLUGIN_SERVER_IDLE: boolean + RETRY_QUEUES: string } export interface PluginsServer extends PluginsServerConfig { @@ -90,21 +89,36 @@ export interface PluginsServer extends PluginsServerConfig { pluginSchedule: Record | null pluginSchedulePromises: Record | null>> eventsProcessor: EventsProcessor + retryQueueManager: RetryQueue } export interface Pausable { - pause: () => Promise - resume: () => void + pause: () => Promise | void + resume: () => Promise | void isPaused: () => boolean } export interface Queue extends Pausable { - start: () => Promise - stop: () => Promise + start: () => Promise | void + stop: () => Promise | void } -export interface Queue { - stop: () => Promise +export type OnRetryCallback = (queue: EnqueuedRetry[]) => Promise | void +export interface EnqueuedRetry { + type: string + payload: Record + timestamp: number + pluginConfigId: number + pluginConfigTeam: number +} + +export interface RetryQueue { + startConsumer: (onRetry: OnRetryCallback) => Promise | void + stopConsumer: () => Promise | void + pauseConsumer: () => Promise | void + resumeConsumer: () => Promise | void + isConsumerPaused: () => boolean + enqueue: (retry: EnqueuedRetry) => Promise | void } export type PluginId = number @@ -186,6 +200,7 @@ export interface PluginConfigVMReponse { teardownPlugin: () => Promise processEvent: (event: PluginEvent) => Promise processEventBatch: (batch: PluginEvent[]) => Promise + onRetry: (task: string, payload: Record) => Promise } tasks: Record } @@ -381,4 +396,9 @@ export interface ScheduleControl { reloadSchedule: () => Promise } +export interface RetryQueueConsumerControl { + stop: () => Promise + resume: () => Promise | void +} + export type IngestEventResponse = { success?: boolean; error?: string } diff --git a/src/worker/plugins/run.ts b/src/worker/plugins/run.ts index 78b3bb59..abdbbbd4 100644 --- a/src/worker/plugins/run.ts +++ b/src/worker/plugins/run.ts @@ -1,7 +1,7 @@ import { PluginEvent } from '@posthog/plugin-scaffold' import { processError } from '../../shared/error' -import { PluginConfig, PluginsServer } from '../../types' +import { EnqueuedRetry, PluginConfig, PluginsServer } from '../../types' export async function runPlugins(server: PluginsServer, event: PluginEvent): Promise { const pluginsToRun = getPluginsForTeam(server, event.team_id) @@ -90,3 +90,18 @@ export async function runPluginTask(server: PluginsServer, taskName: string, plu 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 + try { + const pluginConfig = server.pluginConfigs.get(retry.pluginConfigId) + const task = await pluginConfig?.vm?.getOnRetry() + response = await task?.(retry.type, retry.payload) + } catch (error) { + await processError(server, retry.pluginConfigId, error) + server.statsd?.increment(`plugin.retry.${retry.type}.${retry.pluginConfigId}.ERROR`) + } + server.statsd?.timing(`plugin.retry.${retry.type}.${retry.pluginConfigId}`, timer) + return response +} diff --git a/src/worker/vm/extensions/retry.ts b/src/worker/vm/extensions/retry.ts new file mode 100644 index 00000000..9d4899f3 --- /dev/null +++ b/src/worker/vm/extensions/retry.ts @@ -0,0 +1,23 @@ +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 954a914a..97e0f21d 100644 --- a/src/worker/vm/lazy.ts +++ b/src/worker/vm/lazy.ts @@ -45,6 +45,10 @@ export class LazyPluginVM { return (await this.resolveInternalVm)?.methods.teardownPlugin || null } + async getOnRetry(): Promise { + return (await this.resolveInternalVm)?.methods.onRetry || null + } + async getTask(name: string): Promise { return (await this.resolveInternalVm)?.tasks[name] || null } diff --git a/src/worker/vm/vm.ts b/src/worker/vm/vm.ts index dac08906..423e3b5a 100644 --- a/src/worker/vm/vm.ts +++ b/src/worker/vm/vm.ts @@ -7,6 +7,7 @@ import { createConsole } from './extensions/console' import { createGeoIp } from './extensions/geoip' import { createGoogle } from './extensions/google' import { createPosthog } from './extensions/posthog' +import { createRetry } from './extensions/retry' import { createStorage } from './extensions/storage' import { imports } from './imports' import { transformCode } from './transforms' @@ -65,6 +66,7 @@ export async function createPluginConfigVM( attachments: pluginConfig.attachments, storage: createStorage(server, pluginConfig), geoip: createGeoIp(server), + retry: createRetry(server, pluginConfig), }, '__pluginHostMeta' ) @@ -137,6 +139,7 @@ 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 diff --git a/src/worker/worker.ts b/src/worker/worker.ts index cc06cc79..a16c9748 100644 --- a/src/worker/worker.ts +++ b/src/worker/worker.ts @@ -6,7 +6,7 @@ import { status } from '../shared/status' import { cloneObject } from '../shared/utils' import { PluginsServer, PluginsServerConfig } from '../types' import { ingestEvent } from './ingestion/ingest-event' -import { runPlugins, runPluginsOnBatch, runPluginTask } from './plugins/run' +import { runOnRetry,runPlugins, runPluginsOnBatch, runPluginTask } from './plugins/run' import { loadSchedule, setupPlugins } from './plugins/setup' import { teardownPlugins } from './plugins/teardown' @@ -46,6 +46,9 @@ export const createTaskRunner = (server: PluginsServer): TaskWorker => async ({ // 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) } diff --git a/tests/retry.test.ts b/tests/retry.test.ts new file mode 100644 index 00000000..595517b9 --- /dev/null +++ b/tests/retry.test.ts @@ -0,0 +1,53 @@ +import * as fetch from 'node-fetch' + +import { startPluginsServer } from '../src/main/pluginsServer' +import { delay } from '../src/shared/utils' +import { LogLevel } from '../src/types' +import { makePiscina } from '../src/worker/piscina' +import { createPosthog } from '../src/worker/vm/extensions/posthog' +import { pluginConfig39 } from './helpers/plugins' +import { resetTestDatabase } from './helpers/sql' + +jest.mock('../src/shared/sql') +jest.setTimeout(60000) // 60 sec timeout + +describe('retry queues', () => { + test('on retry gets called', async () => { + const testCode = ` + import fetch from 'node-fetch' + + export async function onRetry (type, payload, meta) { + if (type === 'processEvent') { + console.log('retrying event!', type) + } + void fetch('https://google.com/retry.json?query=' + type) + } + export async function processEvent (event, meta) { + if (event.properties?.hi === 'ha') { + 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) + + expect(fetch).not.toHaveBeenCalled() + + posthog.capture('my event', { hi: 'ha' }) + await delay(5000) + + // can't use this as the call is in a different thread + // expect(fetch).toHaveBeenCalledWith('https://google.com/retry.json?query=processEvent') + + await server.stop() + }) +}) From 9d75e8b05bfbd8d5177e88e0d583d447f199ca25 Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Tue, 20 Apr 2021 21:48:30 +0200 Subject: [PATCH 03/22] capture console.log in tests via a temp file --- .gitignore | 1 + src/main/retry/fs-queue.ts | 19 +++++++--------- src/worker/vm/extensions/test-utils.ts | 31 ++++++++++++++++++++++++++ src/worker/vm/imports.ts | 7 ++++++ tests/retry.test.ts | 24 ++++++++++---------- 5 files changed, 59 insertions(+), 23 deletions(-) create mode 100644 src/worker/vm/extensions/test-utils.ts diff --git a/.gitignore b/.gitignore index 6dce2356..4e1df3f1 100644 --- a/.gitignore +++ b/.gitignore @@ -6,3 +6,4 @@ dist/ yalc.lock .yalc/ src/idl/protos.* +tmp diff --git a/src/main/retry/fs-queue.ts b/src/main/retry/fs-queue.ts index c1de91e0..c6911d80 100644 --- a/src/main/retry/fs-queue.ts +++ b/src/main/retry/fs-queue.ts @@ -1,6 +1,7 @@ import { EnqueuedRetry, OnRetryCallback, RetryQueue } from '../../types' import Timeout = NodeJS.Timeout import * as fs from 'fs' +import * as path from 'path' export class FsQueue implements RetryQueue { paused: boolean @@ -8,15 +9,17 @@ export class FsQueue implements RetryQueue { interval: Timeout | null filename: string - constructor(filename = '/tmp/fs-queue-file.txt') { + constructor(filename?: string) { if (process.env.NODE_ENV !== 'test') { throw new Error('Can not use FsQueue outside tests') } this.paused = false this.started = false this.interval = null - this.filename = filename - // this.queued = [] + this.filename = filename || path.join(process.cwd(), 'tmp', 'fs-queue.txt') + + fs.mkdirSync(path.dirname(this.filename), { recursive: true }) + fs.writeFileSync(this.filename, '') } enqueue(retry: EnqueuedRetry): Promise | void { @@ -35,14 +38,8 @@ export class FsQueue implements RetryQueue { .readFileSync(this.filename) .toString() .split('\n') - .map((s) => { - try { - return JSON.parse(s) as EnqueuedRetry - } catch (e) { - return null - } - }) - .filter((a) => !!a) as EnqueuedRetry[] + .filter((a) => a) + .map((s) => JSON.parse(s) as EnqueuedRetry) const newQueue = queue.filter((element) => element.timestamp < timestamp) if (newQueue.length > 0) { diff --git a/src/worker/vm/extensions/test-utils.ts b/src/worker/vm/extensions/test-utils.ts new file mode 100644 index 00000000..3c0192b7 --- /dev/null +++ b/src/worker/vm/extensions/test-utils.ts @@ -0,0 +1,31 @@ +import fs from 'fs' +import path from 'path' + +const consoleFile = path.join(process.cwd(), 'tmp', 'test-console.txt') + +export const writeToFile = { + console: { + log: (...args: any[]) => { + fs.appendFileSync(consoleFile, `${JSON.stringify(args)}\n`) + }, + reset(): void { + fs.mkdirSync(path.join(process.cwd(), 'tmp'), { recursive: true }) + fs.writeFileSync(consoleFile, '') + }, + read(): any[] { + try { + return fs + .readFileSync(consoleFile) + .toString() + .split('\n') + .filter((str) => !!str) + .map((part) => JSON.parse(part)) + } catch (error) { + if (error.code === 'ENOENT') { + return [] + } + throw error + } + }, + }, +} diff --git a/src/worker/vm/imports.ts b/src/worker/vm/imports.ts index aba2c146..b28a59ef 100644 --- a/src/worker/vm/imports.ts +++ b/src/worker/vm/imports.ts @@ -7,6 +7,8 @@ import fetch from 'node-fetch' import snowflake from 'snowflake-sdk' import * as zlib from 'zlib' +import { writeToFile } from './extensions/test-utils' + export const imports = { crypto: crypto, zlib: zlib, @@ -16,4 +18,9 @@ export const imports = { '@google-cloud/bigquery': { BigQuery }, '@posthog/plugin-contrib': contrib, 'aws-sdk': AWS, + ...(process.env.NODE_ENV === 'test' + ? { + 'test-utils/write-to-file': writeToFile, + } + : {}), } diff --git a/tests/retry.test.ts b/tests/retry.test.ts index 595517b9..bf96c127 100644 --- a/tests/retry.test.ts +++ b/tests/retry.test.ts @@ -1,29 +1,31 @@ -import * as fetch from 'node-fetch' - import { startPluginsServer } from '../src/main/pluginsServer' import { delay } from '../src/shared/utils' import { LogLevel } from '../src/types' import { makePiscina } from '../src/worker/piscina' import { createPosthog } from '../src/worker/vm/extensions/posthog' +import { imports } from '../src/worker/vm/imports' import { pluginConfig39 } from './helpers/plugins' -import { resetTestDatabase } from './helpers/sql' +import { getErrorForPluginConfig, resetTestDatabase } from './helpers/sql' jest.mock('../src/shared/sql') jest.setTimeout(60000) // 60 sec timeout +const { console: testConsole } = imports['test-utils/write-to-file'] + describe('retry queues', () => { + beforeEach(() => { + testConsole.reset() + }) test('on retry gets called', async () => { const testCode = ` - import fetch from 'node-fetch' + import { console } from 'test-utils/write-to-file' export async function onRetry (type, payload, meta) { - if (type === 'processEvent') { - console.log('retrying event!', type) - } - void fetch('https://google.com/retry.json?query=' + type) + 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 @@ -38,15 +40,13 @@ describe('retry queues', () => { }, makePiscina ) + console.log(await getErrorForPluginConfig(39)) const posthog = createPosthog(server.server, pluginConfig39) - expect(fetch).not.toHaveBeenCalled() - posthog.capture('my event', { hi: 'ha' }) await delay(5000) - // can't use this as the call is in a different thread - // expect(fetch).toHaveBeenCalledWith('https://google.com/retry.json?query=processEvent') + expect(testConsole.read()).toEqual([['processEvent'], ['retrying event!', 'processEvent']]) await server.stop() }) From 6b3a8c0d73326b011186cfe06d01304b504907de Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Tue, 20 Apr 2021 22:51:37 +0200 Subject: [PATCH 04/22] add graphile queue --- package.json | 1 + src/main/retry/fs-queue.ts | 4 + src/main/retry/graphile-queue.ts | 88 +++++++++++++++++++ src/main/retry/retry-queue-manager.ts | 8 +- src/main/services/schedule.ts | 6 +- src/shared/server.ts | 1 + src/types.ts | 1 + tests/retry.test.ts | 54 +++++++++++- yarn.lock | 119 +++++++++++++++++++++++++- 9 files changed, 270 insertions(+), 12 deletions(-) create mode 100644 src/main/retry/graphile-queue.ts diff --git a/package.json b/package.json index b7bf2c9a..b7845d3a 100644 --- a/package.json +++ b/package.json @@ -62,6 +62,7 @@ "fast-deep-equal": "^3.1.3", "fastify": "^3.8.0", "generic-pool": "^3.7.1", + "graphile-worker": "^0.11.0", "hot-shots": "^8.2.1", "ioredis": "^4.19.2", "kafkajs": "^1.15.0", diff --git a/src/main/retry/fs-queue.ts b/src/main/retry/fs-queue.ts index c6911d80..49c877ca 100644 --- a/src/main/retry/fs-queue.ts +++ b/src/main/retry/fs-queue.ts @@ -26,6 +26,10 @@ export class FsQueue implements RetryQueue { fs.appendFileSync(this.filename, `${JSON.stringify(retry)}\n`) } + quit(): void { + // nothing to do + } + startConsumer(onRetry: OnRetryCallback): void { fs.writeFileSync(this.filename, '') this.started = true diff --git a/src/main/retry/graphile-queue.ts b/src/main/retry/graphile-queue.ts new file mode 100644 index 00000000..e1cbbc53 --- /dev/null +++ b/src/main/retry/graphile-queue.ts @@ -0,0 +1,88 @@ +import { makeWorkerUtils, run, Runner, WorkerUtils } from 'graphile-worker' + +import { EnqueuedRetry, OnRetryCallback, PluginsServer, RetryQueue } from '../../types' + +export class GraphileQueue implements RetryQueue { + pluginsServer: PluginsServer + started: boolean + paused: boolean + onRetry: OnRetryCallback | null + runner: Runner | null + workerUtils: WorkerUtils | null + + constructor(pluginsServer: PluginsServer) { + this.pluginsServer = pluginsServer + this.started = false + this.paused = false + this.onRetry = null + this.runner = null + this.workerUtils = null + } + + async enqueue(retry: EnqueuedRetry): Promise { + if (!this.workerUtils) { + this.workerUtils = await makeWorkerUtils({ + connectionString: this.pluginsServer.DATABASE_URL, + }) + await this.workerUtils.migrate() + } + await this.workerUtils.addJob('retryTask', retry, { runAt: new Date(retry.timestamp), maxAttempts: 1 }) + } + + async quit(): Promise { + const oldWorkerUtils = this.workerUtils + this.workerUtils = null + await oldWorkerUtils?.release() + } + + async startConsumer(onRetry: OnRetryCallback): Promise { + this.started = true + this.onRetry = onRetry + await this.syncState() + } + + async stopConsumer(): Promise { + this.started = false + await this.syncState() + } + + async pauseConsumer(): Promise { + this.paused = true + await this.syncState() + } + + isConsumerPaused(): boolean { + return this.paused + } + + async resumeConsumer(): Promise { + this.paused = false + await this.syncState() + } + + async syncState(): Promise { + if (this.started && !this.paused) { + if (!this.runner) { + this.runner = await run({ + connectionString: this.pluginsServer.DATABASE_URL, + 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]) + }, + }, + }) + } + } else { + if (this.runner) { + const oldRunner = this.runner + this.runner = null + await oldRunner?.stop() + } + } + } +} diff --git a/src/main/retry/retry-queue-manager.ts b/src/main/retry/retry-queue-manager.ts index ced7635a..60f0bfbe 100644 --- a/src/main/retry/retry-queue-manager.ts +++ b/src/main/retry/retry-queue-manager.ts @@ -2,9 +2,11 @@ import * as Sentry from '@sentry/node' import { EnqueuedRetry, OnRetryCallback, PluginsServer, RetryQueue } from '../../types' import { FsQueue } from './fs-queue' +import { GraphileQueue } from './graphile-queue' const queues = { fs: () => new FsQueue(), + graphile: (pluginsServer: PluginsServer) => new GraphileQueue(pluginsServer), } export class RetryQueueManager implements RetryQueue { @@ -19,7 +21,7 @@ export class RetryQueueManager implements RetryQueue { .map( (queue): RetryQueue => { if (queues[queue as keyof typeof queues]) { - return queues[queue as keyof typeof queues]() + return queues[queue as keyof typeof queues](pluginsServer) } else { throw new Error(`Unknown retry queue "${queue}"`) } @@ -46,6 +48,10 @@ export class RetryQueueManager implements RetryQueue { throw new Error('No RetryQueue available') } + async quit(): Promise { + await Promise.all(this.retryQueues.map((r) => r.quit())) + } + async startConsumer(onRetry: OnRetryCallback): Promise { await Promise.all(this.retryQueues.map((r) => r.startConsumer(onRetry))) } diff --git a/src/main/services/schedule.ts b/src/main/services/schedule.ts index 6514e02f..7131cc4e 100644 --- a/src/main/services/schedule.ts +++ b/src/main/services/schedule.ts @@ -9,11 +9,7 @@ import { startRedlock } from './redlock' export const LOCKED_RESOURCE = 'plugin-server:locks:schedule' -export async function startSchedule( - server: PluginsServer, - piscina: Piscina, - onLock?: () => void -): Promise { +export async function startSchedule(server: PluginsServer, piscina: Piscina): Promise { status.info('⏰', 'Starting scheduling service...') let stopped = false diff --git a/src/shared/server.ts b/src/shared/server.ts index 18cab2db..6e113e11 100644 --- a/src/shared/server.ts +++ b/src/shared/server.ts @@ -164,6 +164,7 @@ export async function createServer( const closeServer = async () => { server.mmdbUpdateJob?.cancel() + await server.retryQueueManager?.quit() if (kafkaProducer) { clearInterval(kafkaProducer.flushInterval) await kafkaProducer.flush() diff --git a/src/types.ts b/src/types.ts index 4b5bb5c1..5dd48571 100644 --- a/src/types.ts +++ b/src/types.ts @@ -119,6 +119,7 @@ export interface RetryQueue { resumeConsumer: () => Promise | void isConsumerPaused: () => boolean enqueue: (retry: EnqueuedRetry) => Promise | void + quit: () => Promise | void } export type PluginId = number diff --git a/tests/retry.test.ts b/tests/retry.test.ts index bf96c127..0520d402 100644 --- a/tests/retry.test.ts +++ b/tests/retry.test.ts @@ -1,11 +1,14 @@ +import { makeWorkerUtils, WorkerUtils } from 'graphile-worker' + import { startPluginsServer } from '../src/main/pluginsServer' +import { getDefaultConfig } from '../src/shared/config' import { delay } from '../src/shared/utils' import { LogLevel } from '../src/types' import { makePiscina } from '../src/worker/piscina' import { createPosthog } from '../src/worker/vm/extensions/posthog' import { imports } from '../src/worker/vm/imports' import { pluginConfig39 } from './helpers/plugins' -import { getErrorForPluginConfig, resetTestDatabase } from './helpers/sql' +import { resetTestDatabase } from './helpers/sql' jest.mock('../src/shared/sql') jest.setTimeout(60000) // 60 sec timeout @@ -13,10 +16,20 @@ jest.setTimeout(60000) // 60 sec timeout const { console: testConsole } = imports['test-utils/write-to-file'] describe('retry queues', () => { - beforeEach(() => { + let workerUtils: WorkerUtils + + beforeEach(async () => { testConsole.reset() + workerUtils = await makeWorkerUtils({ + connectionString: getDefaultConfig().DATABASE_URL, + }) }) - test('on retry gets called', async () => { + + afterEach(async () => { + await workerUtils.release() + }) + + test('onRetry gets called', async () => { const testCode = ` import { console } from 'test-utils/write-to-file' @@ -40,7 +53,40 @@ describe('retry queues', () => { }, makePiscina ) - console.log(await getErrorForPluginConfig(39)) + 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 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' }) diff --git a/yarn.lock b/yarn.lock index ca8f5eec..5fc7f43d 100644 --- a/yarn.lock +++ b/yarn.lock @@ -1200,6 +1200,11 @@ resolved "https://registry.yarnpkg.com/@google-cloud/promisify/-/promisify-2.0.3.tgz#f934b5cdc939e3c7039ff62b9caaf59a9d89e3a8" integrity sha512-d4VSA86eL/AFTe5xtyZX+ePUjE8dIFu2T8zmdeNBSa5/kNgXPCx/o/wbFNHAGLJdGnk1vddRuMESD9HbOC8irw== +"@graphile/logger@^0.2.0": + version "0.2.0" + resolved "https://registry.yarnpkg.com/@graphile/logger/-/logger-0.2.0.tgz#e484ec420162157c6e6f0cfb080fa29ef3a714ba" + integrity sha512-jjcWBokl9eb1gVJ85QmoaQ73CQ52xAaOCF29ukRbYNl6lY+ts0ErTaDYOBlejcbUs2OpaiqYLO5uDhyLFzWw4w== + "@istanbuljs/load-nyc-config@^1.0.0": version "1.1.0" resolved "https://registry.yarnpkg.com/@istanbuljs/load-nyc-config/-/load-nyc-config-1.1.0.tgz#fd3db1d59ecf7cf121e80650bb86712f9b55eced" @@ -1660,6 +1665,11 @@ resolved "https://registry.yarnpkg.com/@types/cookiejar/-/cookiejar-2.1.2.tgz#66ad9331f63fe8a3d3d9d8c6e3906dd10f6446e8" integrity sha512-t73xJJrvdTjXrn4jLS9VSGRbz0nUY3cl2DMGDU48lKl+HR9dbbjW2A9r3g40VA++mQpy6uuHg33gy7du2BKpog== +"@types/debug@^4.1.2": + version "4.1.5" + resolved "https://registry.yarnpkg.com/@types/debug/-/debug-4.1.5.tgz#b14efa8852b7768d898906613c23f688713e02cd" + integrity sha512-Q1y515GcOdTHgagaVFhHnIFQ38ygs/kmxdNpvpou+raI9UO3YZcHDngBSYKQklcKlvA7iuQlmIKbzvmxcOE9CQ== + "@types/generic-pool@^3.1.9": version "3.1.9" resolved "https://registry.yarnpkg.com/@types/generic-pool/-/generic-pool-3.1.9.tgz#cc82ee0d92561fce713f8f9a7b2380eda8a89dcb" @@ -1773,6 +1783,15 @@ resolved "https://registry.yarnpkg.com/@types/pg-types/-/pg-types-1.11.5.tgz#1eebbe62b6772fcc75c18957a90f933d155e005b" integrity sha512-L8ogeT6vDzT1vxlW3KITTCt+BVXXVkLXfZ/XNm6UqbcJgxf+KPO7yjWx7dQQE8RW07KopL10x2gNMs41+IkMGQ== +"@types/pg@^7.14.3": + version "7.14.11" + resolved "https://registry.yarnpkg.com/@types/pg/-/pg-7.14.11.tgz#daf5555504a1f7af4263df265d91f140fece52e3" + integrity sha512-EnZkZ1OMw9DvNfQkn2MTJrwKmhJYDEs5ujWrPfvseWNoI95N8B4HzU/Ltrq5ZfYxDX/Zg8mTzwr6UAyTjjFvXA== + dependencies: + "@types/node" "*" + pg-protocol "^1.2.0" + pg-types "^2.2.0" + "@types/pg@^7.14.6": version "7.14.7" resolved "https://registry.yarnpkg.com/@types/pg/-/pg-7.14.7.tgz#b25532a424f58e70432ac31c77507dfb7b9349a8" @@ -2757,6 +2776,15 @@ cliui@^6.0.0: strip-ansi "^6.0.0" wrap-ansi "^6.2.0" +cliui@^7.0.2: + version "7.0.4" + resolved "https://registry.yarnpkg.com/cliui/-/cliui-7.0.4.tgz#a0265ee655476fc807aea9df3df8df7783808b4f" + integrity sha512-OcRE68cOsVMXp1Yvonl/fzkQOyjLSu/8bhPDfQt0e0/Eb283TKP20Fs2MqoPsr9SwA595rRCA+QMzYc9nBP+JQ== + dependencies: + string-width "^4.2.0" + strip-ansi "^6.0.0" + wrap-ansi "^7.0.0" + cluster-key-slot@^1.1.0: version "1.1.0" resolved "https://registry.yarnpkg.com/cluster-key-slot/-/cluster-key-slot-1.1.0.tgz#30474b2a981fb12172695833052bc0d01336d10d" @@ -3953,7 +3981,7 @@ get-caller-file@^1.0.2: resolved "https://registry.yarnpkg.com/get-caller-file/-/get-caller-file-1.0.3.tgz#f978fa4c90d1dfe7ff2d6beda2a515e713bdcf4a" integrity sha512-3t6rVToeoZfYSGd8YoLFR2DJkiQrIiUrGcjvFX2mDw3bn6k2OtwHN0TNCLbBO+w8qTvimhDkv+LSscbJY1vE6w== -get-caller-file@^2.0.1: +get-caller-file@^2.0.1, get-caller-file@^2.0.5: version "2.0.5" resolved "https://registry.yarnpkg.com/get-caller-file/-/get-caller-file-2.0.5.tgz#4f94412a82db32f36e3b0b9741f8a97feb031f7e" integrity sha512-DyFP3BM/3YHTQOCUL/w0OZHR0lpKeGrxotcHWcqNEdnltqFwXVfhEBQ94eIo34AfQpo0rGki4cyIiftY06h2Fg== @@ -4086,6 +4114,21 @@ graceful-fs@^4.1.11, graceful-fs@^4.1.2, graceful-fs@^4.2.4: resolved "https://registry.yarnpkg.com/graceful-fs/-/graceful-fs-4.2.4.tgz#2256bde14d3632958c465ebc96dc467ca07a29fb" integrity sha512-WjKPNJF79dtJAVniUlGGWHYGz2jWxT6VhN/4m1NdkbZ2nOsEF+cI1Edgql5zCRhs/VsQYRvrXctxktVXZUkixw== +graphile-worker@^0.11.0: + version "0.11.0" + resolved "https://registry.yarnpkg.com/graphile-worker/-/graphile-worker-0.11.0.tgz#eb9984440ebff76f0dd1d3d61c03664064ca255d" + integrity sha512-t3GHGQnafZEWNgBoPaAs0FcR4n4V9uGIY+OlmgBtVl+3zHwC80yBcjW0/yMuCGhOvdiC7Yw9MYlPXcodlypGYA== + dependencies: + "@graphile/logger" "^0.2.0" + "@types/debug" "^4.1.2" + "@types/pg" "^7.14.3" + chokidar "^3.4.0" + cosmiconfig "^7.0.0" + json5 "^2.1.3" + pg ">=6.5 <9" + tslib "^2.1.0" + yargs "^16.2.0" + growly@^1.3.0: version "1.3.0" resolved "https://registry.yarnpkg.com/growly/-/growly-1.3.0.tgz#f10748cbe76af964b7c96c93c6bcc28af120c081" @@ -5185,6 +5228,13 @@ json5@^1.0.1: dependencies: minimist "^1.2.0" +json5@^2.1.3: + version "2.2.0" + resolved "https://registry.yarnpkg.com/json5/-/json5-2.2.0.tgz#2dfefe720c6ba525d9ebd909950f0515316c89a3" + integrity sha512-f+8cldu7X/y7RAJurMEJmdoKXGB/X550w2Nr3tTbezL6RwEE/iMcm+tZnXeoZtKuOq6ft8+CqzEkrIgx1fPoQA== + dependencies: + minimist "^1.2.5" + jsonwebtoken@^8.5.1: version "8.5.1" resolved "https://registry.yarnpkg.com/jsonwebtoken/-/jsonwebtoken-8.5.1.tgz#00e71e0b8df54c2121a1f26137df2280673bcc0d" @@ -6249,6 +6299,11 @@ pg-connection-string@^2.4.0: resolved "https://registry.yarnpkg.com/pg-connection-string/-/pg-connection-string-2.4.0.tgz#c979922eb47832999a204da5dbe1ebf2341b6a10" integrity sha512-3iBXuv7XKvxeMrIgym7njT+HlZkwZqqGX4Bu9cci8xHZNT+Um1gWKqCsAzcC0d95rcKMU5WBg6YRUcHyV0HZKQ== +pg-connection-string@^2.5.0: + version "2.5.0" + resolved "https://registry.yarnpkg.com/pg-connection-string/-/pg-connection-string-2.5.0.tgz#538cadd0f7e603fc09a12590f3b8a452c2c0cf34" + integrity sha512-r5o/V/ORTA6TmUnyWZR9nCj1klXCO2CEKNRlVuJptZe85QuhFayC7WeMic7ndayT5IRIR0S0xFxFi2ousartlQ== + pg-int8@1.0.1: version "1.0.1" resolved "https://registry.yarnpkg.com/pg-int8/-/pg-int8-1.0.1.tgz#943bd463bf5b71b4170115f80f8efc9a0c0eb78c" @@ -6259,12 +6314,22 @@ pg-pool@^3.2.2: resolved "https://registry.yarnpkg.com/pg-pool/-/pg-pool-3.2.2.tgz#a560e433443ed4ad946b84d774b3f22452694dff" integrity sha512-ORJoFxAlmmros8igi608iVEbQNNZlp89diFVx6yV5v+ehmpMY9sK6QgpmgoXbmkNaBAx8cOOZh9g80kJv1ooyA== +pg-pool@^3.3.0: + version "3.3.0" + resolved "https://registry.yarnpkg.com/pg-pool/-/pg-pool-3.3.0.tgz#12d5c7f65ea18a6e99ca9811bd18129071e562fc" + integrity sha512-0O5huCql8/D6PIRFAlmccjphLYWC+JIzvUhSzXSpGaf+tjTZc4nn+Lr7mLXBbFJfvwbP0ywDv73EiaBsxn7zdg== + +pg-protocol@^1.2.0, pg-protocol@^1.5.0: + version "1.5.0" + resolved "https://registry.yarnpkg.com/pg-protocol/-/pg-protocol-1.5.0.tgz#b5dd452257314565e2d54ab3c132adc46565a6a0" + integrity sha512-muRttij7H8TqRNu/DxrAJQITO4Ac7RmX3Klyr/9mJEOBeIpgnF8f9jAfRz5d3XwQZl5qBjF9gLsUtMPJE0vezQ== + pg-protocol@^1.4.0: version "1.4.0" resolved "https://registry.yarnpkg.com/pg-protocol/-/pg-protocol-1.4.0.tgz#43a71a92f6fe3ac559952555aa3335c8cb4908be" integrity sha512-El+aXWcwG/8wuFICMQjM5ZSAm6OWiJicFdNYo+VY3QP+8vI4SvLIWVe51PppTzMhikUJR+PsyIFKqfdXPz/yxA== -pg-types@^2.1.0: +pg-types@^2.1.0, pg-types@^2.2.0: version "2.2.0" resolved "https://registry.yarnpkg.com/pg-types/-/pg-types-2.2.0.tgz#2d0250d636454f7cfa3b6ae0382fdfa8063254a3" integrity sha512-qTAAlrEsl8s4OiEQY69wDvcMIdQN6wdz5ojQiOy6YRMuynxenON0O5oCpJI6lshc6scgAY8qvJ2On/p+CXY0GA== @@ -6275,6 +6340,19 @@ pg-types@^2.1.0: postgres-date "~1.0.4" postgres-interval "^1.1.0" +"pg@>=6.5 <9": + version "8.6.0" + resolved "https://registry.yarnpkg.com/pg/-/pg-8.6.0.tgz#e222296b0b079b280cce106ea991703335487db2" + integrity sha512-qNS9u61lqljTDFvmk/N66EeGq3n6Ujzj0FFyNMGQr6XuEv4tgNTXvJQTfJdcvGit5p5/DWPu+wj920hAJFI+QQ== + dependencies: + buffer-writer "2.0.0" + packet-reader "1.0.0" + pg-connection-string "^2.5.0" + pg-pool "^3.3.0" + pg-protocol "^1.5.0" + pg-types "^2.1.0" + pgpass "1.x" + pg@^8.4.2: version "8.5.1" resolved "https://registry.yarnpkg.com/pg/-/pg-8.5.1.tgz#34dcb15f6db4a29c702bf5031ef2e1e25a06a120" @@ -7771,6 +7849,11 @@ tslib@^1.8.1, tslib@^1.9.0, tslib@^1.9.3: resolved "https://registry.yarnpkg.com/tslib/-/tslib-1.14.1.tgz#cf2d38bdc34a134bcaf1091c41f6619e2f672d00" integrity sha512-Xni35NKzjgMrwevysHTCArtLDpPvye8zV/0E4EyYn43P7/7qvQwPh9BGkHewbMulVntbigmcT7rdX3BNo9wRJg== +tslib@^2.1.0: + version "2.2.0" + resolved "https://registry.yarnpkg.com/tslib/-/tslib-2.2.0.tgz#fb2c475977e35e241311ede2693cee1ec6698f5c" + integrity sha512-gS9GVHRU+RGn5KQM2rllAlR3dU6m7AcpJKdtH8gFvQiC4Otgk98XnmMU+nZenHt/+VhnBPWwgrJsyrdcw6i23w== + tsutils@^3.17.1: version "3.17.1" resolved "https://registry.yarnpkg.com/tsutils/-/tsutils-3.17.1.tgz#ed719917f11ca0dee586272b2ac49e015a2dd759" @@ -8089,6 +8172,15 @@ wrap-ansi@^6.2.0: string-width "^4.1.0" strip-ansi "^6.0.0" +wrap-ansi@^7.0.0: + version "7.0.0" + resolved "https://registry.yarnpkg.com/wrap-ansi/-/wrap-ansi-7.0.0.tgz#67e145cff510a6a6984bdf1152911d69d2eb9e43" + integrity sha512-YVGIj2kamLSTxw6NsZjoBxfSwsn0ycdesmc4p+Q21c5zPuZ1pl+NfxVdxPtdHvmNVOQ6XSYG4AUtyt/Fi7D16Q== + dependencies: + ansi-styles "^4.0.0" + string-width "^4.1.0" + strip-ansi "^6.0.0" + wrappy@1: version "1.0.2" resolved "https://registry.yarnpkg.com/wrappy/-/wrappy-1.0.2.tgz#b5243d8f3ec1aa35f1364605bc0d1036e30ab69f" @@ -8142,6 +8234,11 @@ y18n@^4.0.0: resolved "https://registry.yarnpkg.com/y18n/-/y18n-4.0.1.tgz#8db2b83c31c5d75099bb890b23f3094891e247d4" integrity sha512-wNcy4NvjMYL8gogWWYAO7ZFWFfHcbdbE57tZO8e4cbpj8tfUcwrwqSl3ad8HxpYWCdXcJUCeKKZS62Av1affwQ== +y18n@^5.0.5: + version "5.0.8" + resolved "https://registry.yarnpkg.com/y18n/-/y18n-5.0.8.tgz#7f4934d0f7ca8c56f95314939ddcd2dd91ce1d55" + integrity sha512-0pfFzegeDWJHJIAmTLRP2DwHjdF5s7jo9tuztdQxAhINCdvS+3nGINqPd00AphqJR/0LhANUS6/+7SCb98YOfA== + yallist@^4.0.0: version "4.0.0" resolved "https://registry.yarnpkg.com/yallist/-/yallist-4.0.0.tgz#9bb92790d9c0effec63be73519e11a35019a3a72" @@ -8165,6 +8262,11 @@ yargs-parser@^18.1.2: camelcase "^5.0.0" decamelize "^1.2.0" +yargs-parser@^20.2.2: + version "20.2.7" + resolved "https://registry.yarnpkg.com/yargs-parser/-/yargs-parser-20.2.7.tgz#61df85c113edfb5a7a4e36eb8aa60ef423cbc90a" + integrity sha512-FiNkvbeHzB/syOjIUxFDCnhSfzAL8R5vs40MgLFBorXACCOAEaWu0gRZl14vG8MR9AOJIZbmkjhusqBYZ3HTHw== + yargs@^15.4.1: version "15.4.1" resolved "https://registry.yarnpkg.com/yargs/-/yargs-15.4.1.tgz#0d87a16de01aee9d8bec2bfbf74f67851730f4f8" @@ -8182,6 +8284,19 @@ yargs@^15.4.1: y18n "^4.0.0" yargs-parser "^18.1.2" +yargs@^16.2.0: + version "16.2.0" + resolved "https://registry.yarnpkg.com/yargs/-/yargs-16.2.0.tgz#1c82bf0f6b6a66eafce7ef30e376f49a12477f66" + integrity sha512-D1mvvtDG0L5ft/jGWkLpG1+m0eQxOfaBvTNELraWj22wSVUMWxZUvYgJYcKh6jGGIkJFhH4IZPQhR4TKpc8mBw== + dependencies: + cliui "^7.0.2" + escalade "^3.1.1" + get-caller-file "^2.0.5" + require-directory "^2.1.1" + string-width "^4.2.0" + y18n "^5.0.5" + yargs-parser "^20.2.2" + yn@3.1.1: version "3.1.1" resolved "https://registry.yarnpkg.com/yn/-/yn-3.1.1.tgz#1e87401a09d767c1d5eab26a6e4c185182d2eb50" From c75b7a60f38e0e1b7d54d3ecc06a514ff799f901 Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Tue, 20 Apr 2021 22:56:39 +0200 Subject: [PATCH 05/22] make it prettier and safe --- src/main/retry/retry-queue-manager.ts | 1 + src/worker/worker.ts | 2 +- 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/src/main/retry/retry-queue-manager.ts b/src/main/retry/retry-queue-manager.ts index 60f0bfbe..c0b6d831 100644 --- a/src/main/retry/retry-queue-manager.ts +++ b/src/main/retry/retry-queue-manager.ts @@ -18,6 +18,7 @@ export class RetryQueueManager implements RetryQueue { this.retryQueues = pluginsServer.RETRY_QUEUES.split(',') .map((q) => q.trim()) + .filter((q) => !!q) .map( (queue): RetryQueue => { if (queues[queue as keyof typeof queues]) { diff --git a/src/worker/worker.ts b/src/worker/worker.ts index a16c9748..53b0fffb 100644 --- a/src/worker/worker.ts +++ b/src/worker/worker.ts @@ -6,7 +6,7 @@ import { status } from '../shared/status' import { cloneObject } from '../shared/utils' import { PluginsServer, PluginsServerConfig } from '../types' import { ingestEvent } from './ingestion/ingest-event' -import { runOnRetry,runPlugins, runPluginsOnBatch, runPluginTask } from './plugins/run' +import { runOnRetry, runPlugins, runPluginsOnBatch, runPluginTask } from './plugins/run' import { loadSchedule, setupPlugins } from './plugins/setup' import { teardownPlugins } from './plugins/teardown' From 5eab403865a93249846865fdbf7547643898c885 Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Tue, 20 Apr 2021 23:00:42 +0200 Subject: [PATCH 06/22] style fixes --- benchmarks/clickhouse/e2e.kafka.benchmark.ts | 2 +- benchmarks/clickhouse/e2e.timeout.benchmark.ts | 2 +- benchmarks/postgres/e2e.celery.benchmark.ts | 2 +- src/main/queue.ts | 2 +- src/worker/vm/extensions/test-utils.ts | 2 +- tests/postgres/queue.test.ts | 2 +- 6 files changed, 6 insertions(+), 6 deletions(-) diff --git a/benchmarks/clickhouse/e2e.kafka.benchmark.ts b/benchmarks/clickhouse/e2e.kafka.benchmark.ts index 07e3b1db..31880c69 100644 --- a/benchmarks/clickhouse/e2e.kafka.benchmark.ts +++ b/benchmarks/clickhouse/e2e.kafka.benchmark.ts @@ -73,7 +73,7 @@ describe('e2e kafka & clickhouse benchmark', () => { // hope that 5sec is enough to load kafka with all the events (posthog.capture can't be awaited) await delay(5000) - queue.resume() + await queue.resume() console.log('Starting timer') const startTime = performance.now() diff --git a/benchmarks/clickhouse/e2e.timeout.benchmark.ts b/benchmarks/clickhouse/e2e.timeout.benchmark.ts index f8c43ec7..9ecaea35 100644 --- a/benchmarks/clickhouse/e2e.timeout.benchmark.ts +++ b/benchmarks/clickhouse/e2e.timeout.benchmark.ts @@ -71,7 +71,7 @@ describe('e2e kafka processing timeout benchmark', () => { // hope that 5sec is enough to load kafka with all the events (posthog.capture can't be awaited) await delay(5000) - queue.resume() + await queue.resume() console.log('Starting timer') const startTime = performance.now() diff --git a/benchmarks/postgres/e2e.celery.benchmark.ts b/benchmarks/postgres/e2e.celery.benchmark.ts index cf11be2a..0c4d02b5 100644 --- a/benchmarks/postgres/e2e.celery.benchmark.ts +++ b/benchmarks/postgres/e2e.celery.benchmark.ts @@ -73,7 +73,7 @@ describe('e2e celery & postgres benchmark', () => { } await delay(3000) expect(await redis.llen(server.PLUGINS_CELERY_QUEUE)).toEqual(count) - queue.resume() + await queue.resume() console.log('Starting timer') const startTime = performance.now() diff --git a/src/main/queue.ts b/src/main/queue.ts index 29a31720..d135dec5 100644 --- a/src/main/queue.ts +++ b/src/main/queue.ts @@ -18,7 +18,7 @@ export function pauseQueueIfWorkerFull( pause: undefined | (() => void | Promise), server: PluginsServer, piscina?: Piscina -) { +): void { if (pause && (piscina?.queueSize || 0) > (server.WORKER_CONCURRENCY || 4) * (server.WORKER_CONCURRENCY || 4)) { void pause() } diff --git a/src/worker/vm/extensions/test-utils.ts b/src/worker/vm/extensions/test-utils.ts index 3c0192b7..c0a1585f 100644 --- a/src/worker/vm/extensions/test-utils.ts +++ b/src/worker/vm/extensions/test-utils.ts @@ -5,7 +5,7 @@ const consoleFile = path.join(process.cwd(), 'tmp', 'test-console.txt') export const writeToFile = { console: { - log: (...args: any[]) => { + log: (...args: any[]): void => { fs.appendFileSync(consoleFile, `${JSON.stringify(args)}\n`) }, reset(): void { diff --git a/tests/postgres/queue.test.ts b/tests/postgres/queue.test.ts index a28b43c3..2c41f5d7 100644 --- a/tests/postgres/queue.test.ts +++ b/tests/postgres/queue.test.ts @@ -83,7 +83,7 @@ test('pause and resume queue', async () => { expect(await redis.llen(server.PLUGINS_CELERY_QUEUE)).toBe(pluginQueue) - queue.resume() + await queue.resume() await delay(500) From 185a187f62f5da253f7f321ae965c88e21e0b633 Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Tue, 20 Apr 2021 23:21:10 +0200 Subject: [PATCH 07/22] fix some tests --- src/main/services/schedule.ts | 7 ++++++- tests/plugins.test.ts | 9 ++++++++- tests/postgres/vm.test.ts | 1 + 3 files changed, 15 insertions(+), 2 deletions(-) diff --git a/src/main/services/schedule.ts b/src/main/services/schedule.ts index 7131cc4e..8dd563de 100644 --- a/src/main/services/schedule.ts +++ b/src/main/services/schedule.ts @@ -9,7 +9,11 @@ import { startRedlock } from './redlock' export const LOCKED_RESOURCE = 'plugin-server:locks:schedule' -export async function startSchedule(server: PluginsServer, piscina: Piscina): Promise { +export async function startSchedule( + server: PluginsServer, + piscina: Piscina, + onLock?: () => void +): Promise { status.info('⏰', 'Starting scheduling service...') let stopped = false @@ -41,6 +45,7 @@ export async function startSchedule(server: PluginsServer, piscina: Piscina): Pr LOCKED_RESOURCE, () => { weHaveTheLock = true + onLock?.() }, () => { weHaveTheLock = false diff --git a/tests/plugins.test.ts b/tests/plugins.test.ts index 6aad5385..53ba6d73 100644 --- a/tests/plugins.test.ts +++ b/tests/plugins.test.ts @@ -69,7 +69,13 @@ test('setupPlugins and runPlugins', async () => { }) expect(pluginConfig.vm).toBeDefined() const vm = await pluginConfig.vm!.resolveInternalVm - expect(Object.keys(vm!.methods)).toEqual(['setupPlugin', 'teardownPlugin', 'processEvent', 'processEventBatch']) + expect(Object.keys(vm!.methods).sort()).toEqual([ + 'onRetry', + 'processEvent', + 'processEventBatch', + 'setupPlugin', + 'teardownPlugin', + ]) expect(clearError).toHaveBeenCalledWith(mockServer, pluginConfig) @@ -120,6 +126,7 @@ test('plugin meta has what it should have', async () => { 'config', 'geoip', 'global', + 'retry', 'storage', ]) expect(returnedEvent!.properties!['attachments']).toEqual({ diff --git a/tests/postgres/vm.test.ts b/tests/postgres/vm.test.ts index 24411636..f4ff00fd 100644 --- a/tests/postgres/vm.test.ts +++ b/tests/postgres/vm.test.ts @@ -39,6 +39,7 @@ test('empty plugins', async () => { expect(Object.keys(vm).sort()).toEqual(['methods', 'tasks', 'vm']) expect(Object.keys(vm.methods).sort()).toEqual([ + 'onRetry', 'processEvent', 'processEventBatch', 'setupPlugin', From 39c28af14622f1dec4378c134f6f21acdafae0b9 Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Tue, 20 Apr 2021 23:47:24 +0200 Subject: [PATCH 08/22] release if there --- tests/retry.test.ts | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/tests/retry.test.ts b/tests/retry.test.ts index 0520d402..ba2a4fae 100644 --- a/tests/retry.test.ts +++ b/tests/retry.test.ts @@ -16,7 +16,7 @@ jest.setTimeout(60000) // 60 sec timeout const { console: testConsole } = imports['test-utils/write-to-file'] describe('retry queues', () => { - let workerUtils: WorkerUtils + let workerUtils: WorkerUtils | null beforeEach(async () => { testConsole.reset() @@ -26,7 +26,8 @@ describe('retry queues', () => { }) afterEach(async () => { - await workerUtils.release() + await workerUtils?.release() + workerUtils = null }) test('onRetry gets called', async () => { From 9665f44eabcf9f92eb32a40cc4b72318f0e33c88 Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Wed, 21 Apr 2021 09:07:34 +0200 Subject: [PATCH 09/22] split postgres tests --- .github/workflows/ci.yml | 81 ++++++++++++++++++++++++++++++++++++++-- package.json | 3 +- 2 files changed, 80 insertions(+), 4 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 27a070cf..f1f066fb 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -25,8 +25,8 @@ jobs: - name: Lint with ESLint run: yarn lint - tests-postgres: - name: Tests / Postgres + Redis + tests-postgres-1: + name: Tests / Postgres + Redis (1) runs-on: ubuntu-20.04 services: @@ -98,7 +98,82 @@ jobs: # Below DB name has `test_` prepended, as that's how Django (ran above) creates the test DB DATABASE_URL: 'postgres://postgres:postgres@localhost:5432/test_posthog' REDIS_URL: 'redis://localhost' - run: cd plugin-server && yarn test:postgres + run: cd plugin-server && yarn test:postgres:1 + + tests-postgres-2: + name: Tests / Postgres + Redis (2) + runs-on: ubuntu-20.04 + + services: + postgres: + image: postgres:12 + env: + POSTGRES_USER: postgres + POSTGRES_PASSWORD: postgres + POSTGRES_DB: test_posthog + ports: ['5432:5432'] + options: --health-cmd pg_isready --health-interval 10s --health-timeout 5s --health-retries 5 + redis: + image: redis + ports: + - '6379:6379' + options: >- + --health-cmd "redis-cli ping" + --health-interval 10s + --health-timeout 5s + --health-retries 5 + + env: + REDIS_URL: 'redis://localhost' + + steps: + - name: Check out Django server for database setup + uses: actions/checkout@v2 + with: + repository: 'PostHog/posthog' + path: 'posthog/' + + - name: Check out plugin server + uses: actions/checkout@v2 + with: + path: 'plugin-server' + + - name: Set up Python + uses: actions/setup-python@v2 + with: + python-version: 3.8 + + - name: Set up Node 14 + uses: actions/setup-node@v2 + with: + node-version: 14 + + - uses: actions/cache@v2 + with: + path: ${{ env.pythonLocation }} + key: ${{ env.pythonLocation }}-${{ hashFiles('posthog/requirements.txt') }} + + - name: Install requirements.txt dependencies with pip + run: | + pip install --upgrade pip + pip install --upgrade --upgrade-strategy eager -r posthog/requirements.txt + + - name: Set up databases + env: + SECRET_KEY: 'abcdef' # unsafe - for testing only + DATABASE_URL: 'postgres://postgres:postgres@localhost:5432/posthog' + TEST: 'true' + run: python posthog/manage.py setup_test_environment + + - name: Install package.json dependencies with Yarn + run: cd plugin-server && yarn + + - name: Test with Jest + env: + # Below DB name has `test_` prepended, as that's how Django (ran above) creates the test DB + DATABASE_URL: 'postgres://postgres:postgres@localhost:5432/test_posthog' + REDIS_URL: 'redis://localhost' + run: cd plugin-server && yarn test:postgres:2 tests-clickhouse-1: name: Tests / ClickHouse + Kafka (1) diff --git a/package.json b/package.json index b7845d3a..f6aa9a19 100644 --- a/package.json +++ b/package.json @@ -6,7 +6,8 @@ "main": "dist/index.js", "scripts": { "test": "jest --runInBand --forceExit tests/**/*.test.ts", - "test:postgres": "jest --runInBand --forceExit tests/postgres/*.test.ts tests/*.test.ts", + "test:postgres:1": "jest --runInBand --forceExit tests/postgres/*.test.ts", + "test:postgres:2": "jest --runInBand --forceExit tests/*.test.ts", "test:clickhouse:1": "jest --runInBand --forceExit tests/clickhouse/postgres-parity.test.ts tests/clickhouse/e2e.test.ts tests/clickhouse/ingestion-utils.test.ts", "test:clickhouse:2": "jest --runInBand --forceExit tests/clickhouse/process-event.test.ts", "benchmark": "yarn run benchmarks:clickhouse && yarn run benchmark:postgres && yarn run benchmarks:vm", From 6e3e034443c488a16be82de2748f9e9715a77a93 Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Wed, 21 Apr 2021 09:11:40 +0200 Subject: [PATCH 10/22] don't make a graphile worker in all tests --- tests/retry.test.ts | 14 ++------------ 1 file changed, 2 insertions(+), 12 deletions(-) diff --git a/tests/retry.test.ts b/tests/retry.test.ts index ba2a4fae..7c522807 100644 --- a/tests/retry.test.ts +++ b/tests/retry.test.ts @@ -1,7 +1,7 @@ import { makeWorkerUtils, WorkerUtils } from 'graphile-worker' import { startPluginsServer } from '../src/main/pluginsServer' -import { getDefaultConfig } from '../src/shared/config' +import { defaultConfig } from '../src/shared/config' import { delay } from '../src/shared/utils' import { LogLevel } from '../src/types' import { makePiscina } from '../src/worker/piscina' @@ -16,18 +16,8 @@ jest.setTimeout(60000) // 60 sec timeout const { console: testConsole } = imports['test-utils/write-to-file'] describe('retry queues', () => { - let workerUtils: WorkerUtils | null - - beforeEach(async () => { + beforeEach(() => { testConsole.reset() - workerUtils = await makeWorkerUtils({ - connectionString: getDefaultConfig().DATABASE_URL, - }) - }) - - afterEach(async () => { - await workerUtils?.release() - workerUtils = null }) test('onRetry gets called', async () => { From 3e34bb2155aa008d7c72b159a60ab4c43662710f Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Wed, 21 Apr 2021 10:02:06 +0200 Subject: [PATCH 11/22] revert "split postgres tests" --- .github/workflows/ci.yml | 81 ++-------------------------------------- package.json | 3 +- 2 files changed, 4 insertions(+), 80 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index f1f066fb..27a070cf 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -25,8 +25,8 @@ jobs: - name: Lint with ESLint run: yarn lint - tests-postgres-1: - name: Tests / Postgres + Redis (1) + tests-postgres: + name: Tests / Postgres + Redis runs-on: ubuntu-20.04 services: @@ -98,82 +98,7 @@ jobs: # Below DB name has `test_` prepended, as that's how Django (ran above) creates the test DB DATABASE_URL: 'postgres://postgres:postgres@localhost:5432/test_posthog' REDIS_URL: 'redis://localhost' - run: cd plugin-server && yarn test:postgres:1 - - tests-postgres-2: - name: Tests / Postgres + Redis (2) - runs-on: ubuntu-20.04 - - services: - postgres: - image: postgres:12 - env: - POSTGRES_USER: postgres - POSTGRES_PASSWORD: postgres - POSTGRES_DB: test_posthog - ports: ['5432:5432'] - options: --health-cmd pg_isready --health-interval 10s --health-timeout 5s --health-retries 5 - redis: - image: redis - ports: - - '6379:6379' - options: >- - --health-cmd "redis-cli ping" - --health-interval 10s - --health-timeout 5s - --health-retries 5 - - env: - REDIS_URL: 'redis://localhost' - - steps: - - name: Check out Django server for database setup - uses: actions/checkout@v2 - with: - repository: 'PostHog/posthog' - path: 'posthog/' - - - name: Check out plugin server - uses: actions/checkout@v2 - with: - path: 'plugin-server' - - - name: Set up Python - uses: actions/setup-python@v2 - with: - python-version: 3.8 - - - name: Set up Node 14 - uses: actions/setup-node@v2 - with: - node-version: 14 - - - uses: actions/cache@v2 - with: - path: ${{ env.pythonLocation }} - key: ${{ env.pythonLocation }}-${{ hashFiles('posthog/requirements.txt') }} - - - name: Install requirements.txt dependencies with pip - run: | - pip install --upgrade pip - pip install --upgrade --upgrade-strategy eager -r posthog/requirements.txt - - - name: Set up databases - env: - SECRET_KEY: 'abcdef' # unsafe - for testing only - DATABASE_URL: 'postgres://postgres:postgres@localhost:5432/posthog' - TEST: 'true' - run: python posthog/manage.py setup_test_environment - - - name: Install package.json dependencies with Yarn - run: cd plugin-server && yarn - - - name: Test with Jest - env: - # Below DB name has `test_` prepended, as that's how Django (ran above) creates the test DB - DATABASE_URL: 'postgres://postgres:postgres@localhost:5432/test_posthog' - REDIS_URL: 'redis://localhost' - run: cd plugin-server && yarn test:postgres:2 + run: cd plugin-server && yarn test:postgres tests-clickhouse-1: name: Tests / ClickHouse + Kafka (1) diff --git a/package.json b/package.json index f6aa9a19..b7845d3a 100644 --- a/package.json +++ b/package.json @@ -6,8 +6,7 @@ "main": "dist/index.js", "scripts": { "test": "jest --runInBand --forceExit tests/**/*.test.ts", - "test:postgres:1": "jest --runInBand --forceExit tests/postgres/*.test.ts", - "test:postgres:2": "jest --runInBand --forceExit tests/*.test.ts", + "test:postgres": "jest --runInBand --forceExit tests/postgres/*.test.ts tests/*.test.ts", "test:clickhouse:1": "jest --runInBand --forceExit tests/clickhouse/postgres-parity.test.ts tests/clickhouse/e2e.test.ts tests/clickhouse/ingestion-utils.test.ts", "test:clickhouse:2": "jest --runInBand --forceExit tests/clickhouse/process-event.test.ts", "benchmark": "yarn run benchmarks:clickhouse && yarn run benchmark:postgres && yarn run benchmarks:vm", From c4845e6affaf6f73ac6b43b9d1928a2315807615 Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Wed, 28 Apr 2021 12:41:01 +0200 Subject: [PATCH 12/22] skip retries if pluginConfig not found --- src/worker/plugins/run.ts | 22 +++++++++++++++------- 1 file changed, 15 insertions(+), 7 deletions(-) diff --git a/src/worker/plugins/run.ts b/src/worker/plugins/run.ts index 90470042..0d826ce8 100644 --- a/src/worker/plugins/run.ts +++ b/src/worker/plugins/run.ts @@ -1,4 +1,5 @@ import { PluginEvent } from '@posthog/plugin-scaffold' +import * as Sentry from '@sentry/node' import { processError } from '../../shared/error' import { EnqueuedRetry, PluginConfig, PluginsServer } from '../../types' @@ -119,13 +120,20 @@ function getPluginsForTeam(server: PluginsServer, teamId: number): PluginConfig[ export async function runOnRetry(server: PluginsServer, retry: EnqueuedRetry): Promise { const timer = new Date() let response - try { - const pluginConfig = server.pluginConfigs.get(retry.pluginConfigId) - const task = await pluginConfig?.vm?.getOnRetry() - response = await task?.(retry.type, retry.payload) - } catch (error) { - await processError(server, retry.pluginConfigId, error) - server.statsd?.increment(`plugin.retry.${retry.type}.${retry.pluginConfigId}.ERROR`) + 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 From 6fe7eba56545421e05dfe40a9e88985b47aad65a Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Thu, 29 Apr 2021 14:32:05 +0200 Subject: [PATCH 13/22] reset graphile schema before test --- tests/helpers/graphile.ts | 23 +++++++++ tests/retry.test.ts | 98 +++++++++++++++++++++------------------ 2 files changed, 75 insertions(+), 46 deletions(-) create mode 100644 tests/helpers/graphile.ts diff --git a/tests/helpers/graphile.ts b/tests/helpers/graphile.ts new file mode 100644 index 00000000..c0361b02 --- /dev/null +++ b/tests/helpers/graphile.ts @@ -0,0 +1,23 @@ +import { makeWorkerUtils } from 'graphile-worker' +import { Pool } from 'pg' + +import { defaultConfig } from '../../src/shared/config' +import { status } from '../../src/shared/status' + +export async function resetGraphileSchema() { + const db = new Pool({ connectionString: defaultConfig.DATABASE_URL }) + + try { + await db.query('DROP SCHEMA graphile_worker CASCADE') + } catch (e) { + status.error('😱', `Could not dump graphile_worker schema: ${e.message}`) + } finally { + await db.end() + } + + const workerUtils = await makeWorkerUtils({ + connectionString: defaultConfig.DATABASE_URL, + }) + await workerUtils.migrate() + await workerUtils.release() +} diff --git a/tests/retry.test.ts b/tests/retry.test.ts index 7c522807..eafc7228 100644 --- a/tests/retry.test.ts +++ b/tests/retry.test.ts @@ -1,12 +1,10 @@ -import { makeWorkerUtils, WorkerUtils } from 'graphile-worker' - import { startPluginsServer } from '../src/main/pluginsServer' -import { defaultConfig } from '../src/shared/config' import { delay } from '../src/shared/utils' import { LogLevel } from '../src/types' 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' @@ -20,8 +18,9 @@ describe('retry queues', () => { testConsole.reset() }) - test('onRetry gets called', async () => { - const testCode = ` + 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) { @@ -35,56 +34,63 @@ describe('retry queues', () => { 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) + 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(5000) + posthog.capture('my event', { hi: 'ha' }) + await delay(10000) - expect(testConsole.read()).toEqual([['processEvent'], ['retrying event!', 'processEvent']]) + expect(testConsole.read()).toEqual([['processEvent'], ['retrying event!', 'processEvent']]) - await server.stop() + await server.stop() + }) }) - test('graphile retry queue', async () => { - const testCode = ` - import { console } from 'test-utils/write-to-file' + describe('graphile', () => { + beforeEach(async () => { + await resetGraphileSchema() + }) - 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) + 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) } - 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) + 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) + posthog.capture('my event', { hi: 'ha' }) + await delay(5000) - expect(testConsole.read()).toEqual([['processEvent'], ['retrying event!', 'processEvent']]) + expect(testConsole.read()).toEqual([['processEvent'], ['retrying event!', 'processEvent']]) - await server.stop() + await server.stop() + }) }) }) From 16e14b614221e664c9fa2c2fde75ed90a3de656f Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Thu, 29 Apr 2021 14:41:20 +0200 Subject: [PATCH 14/22] fix failing tests by clearing the retry consumer redlock --- tests/retry.test.ts | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/tests/retry.test.ts b/tests/retry.test.ts index eafc7228..984ebd0f 100644 --- a/tests/retry.test.ts +++ b/tests/retry.test.ts @@ -1,4 +1,6 @@ import { startPluginsServer } from '../src/main/pluginsServer' +import { LOCKED_RESOURCE } from '../src/main/services/retry-queue-consumer' +import { createServer } from '../src/shared/server' import { delay } from '../src/shared/utils' import { LogLevel } from '../src/types' import { makePiscina } from '../src/worker/piscina' @@ -14,8 +16,14 @@ jest.setTimeout(60000) // 60 sec timeout const { console: testConsole } = imports['test-utils/write-to-file'] describe('retry queues', () => { - beforeEach(() => { + 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', () => { From c20c67c9215bc96d9be08c289560a1d4bf314eca Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Thu, 29 Apr 2021 14:47:43 +0200 Subject: [PATCH 15/22] bust github actions cache --- .github/workflows/ci.yml | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index f1f066fb..daa75793 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -76,7 +76,7 @@ jobs: - uses: actions/cache@v2 with: path: ${{ env.pythonLocation }} - key: ${{ env.pythonLocation }}-${{ hashFiles('posthog/requirements.txt') }} + key: ${{ env.pythonLocation }}-v1-${{ hashFiles('posthog/requirements.txt') }} - name: Install requirements.txt dependencies with pip run: | @@ -151,7 +151,7 @@ jobs: - uses: actions/cache@v2 with: path: ${{ env.pythonLocation }} - key: ${{ env.pythonLocation }}-${{ hashFiles('posthog/requirements.txt') }} + key: ${{ env.pythonLocation }}-v1-${{ hashFiles('posthog/requirements.txt') }} - name: Install requirements.txt dependencies with pip run: | @@ -241,7 +241,7 @@ jobs: - uses: actions/cache@v2 with: path: ${{ env.pythonLocation }} - key: ${{ env.pythonLocation }}-${{ hashFiles('posthog/requirements.txt') }} + key: ${{ env.pythonLocation }}-v1-${{ hashFiles('posthog/requirements.txt') }} - name: Install requirements.txt dependencies with pip run: | @@ -332,7 +332,7 @@ jobs: - uses: actions/cache@v2 with: path: ${{ env.pythonLocation }} - key: ${{ env.pythonLocation }}-${{ hashFiles('posthog/requirements.txt') }} + key: ${{ env.pythonLocation }}-v1-${{ hashFiles('posthog/requirements.txt') }} - name: Install requirements.txt dependencies with pip run: | @@ -423,7 +423,7 @@ jobs: - uses: actions/cache@v2 with: path: ${{ env.pythonLocation }} - key: ${{ env.pythonLocation }}-${{ hashFiles('posthog/requirements.txt') }} + key: ${{ env.pythonLocation }}-v1-${{ hashFiles('posthog/requirements.txt') }} - name: Install requirements.txt dependencies with pip run: | @@ -499,7 +499,7 @@ jobs: - uses: actions/cache@v2 with: path: ${{ env.pythonLocation }} - key: ${{ env.pythonLocation }}-${{ hashFiles('posthog/requirements.txt') }} + key: ${{ env.pythonLocation }}-v1-${{ hashFiles('posthog/requirements.txt') }} - name: Install requirements.txt dependencies with pip run: | @@ -574,7 +574,7 @@ jobs: - uses: actions/cache@v2 with: path: ${{ env.pythonLocation }} - key: ${{ env.pythonLocation }}-${{ hashFiles('posthog/requirements.txt') }} + key: ${{ env.pythonLocation }}-v1-${{ hashFiles('posthog/requirements.txt') }} - name: Install requirements.txt dependencies with pip run: | @@ -649,7 +649,7 @@ jobs: - uses: actions/cache@v2 with: path: ${{ env.pythonLocation }} - key: ${{ env.pythonLocation }}-${{ hashFiles('posthog/requirements.txt') }} + key: ${{ env.pythonLocation }}-v1-${{ hashFiles('posthog/requirements.txt') }} - name: Install requirements.txt dependencies with pip run: | From a2f32e67d2724eb93d37b99e4ef5b4996b688038 Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Thu, 29 Apr 2021 15:14:18 +0200 Subject: [PATCH 16/22] slight cleanup --- src/shared/config.ts | 2 +- tests/helpers/graphile.ts | 2 +- tests/retry.test.ts | 22 +++++++++++----------- 3 files changed, 13 insertions(+), 13 deletions(-) diff --git a/src/shared/config.ts b/src/shared/config.ts index 8adbb033..b2dba522 100644 --- a/src/shared/config.ts +++ b/src/shared/config.ts @@ -103,7 +103,7 @@ 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: 'queue system and fallbacks to use for retries', + RETRY_QUEUES: 'retry queue engine and fallback queues', STALENESS_RESTART_SECONDS: 'trigger a restart if no event ingested for this duration', } } diff --git a/tests/helpers/graphile.ts b/tests/helpers/graphile.ts index c0361b02..55960f0f 100644 --- a/tests/helpers/graphile.ts +++ b/tests/helpers/graphile.ts @@ -4,7 +4,7 @@ import { Pool } from 'pg' import { defaultConfig } from '../../src/shared/config' import { status } from '../../src/shared/status' -export async function resetGraphileSchema() { +export async function resetGraphileSchema(): Promise { const db = new Pool({ connectionString: defaultConfig.DATABASE_URL }) try { diff --git a/tests/retry.test.ts b/tests/retry.test.ts index 984ebd0f..6ca0d915 100644 --- a/tests/retry.test.ts +++ b/tests/retry.test.ts @@ -29,19 +29,19 @@ describe('retry queues', () => { describe('fs queue', () => { test('onRetry gets called', async () => { const testCode = ` - import { console } from 'test-utils/write-to-file' + 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) + export async function onRetry (type, payload, meta) { + console.log('retrying event!', type) } - return event - } - ` + 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( { From b1ee2c9b7e8b16afa655afd457625ead8ade6218 Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Thu, 29 Apr 2021 15:15:58 +0200 Subject: [PATCH 17/22] fix github/eslint complaining about an `any` --- src/shared/utils.ts | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/src/shared/utils.ts b/src/shared/utils.ts index 774bada7..e07099ca 100644 --- a/src/shared/utils.ts +++ b/src/shared/utils.ts @@ -455,12 +455,11 @@ export enum NodeEnv { Test = 'test', } -export function stringToBoolean(value: any): boolean { +export function stringToBoolean(value: unknown): boolean { if (!value) { return false } - value = String(value) - return ['y', 'yes', 't', 'true', 'on', '1'].includes(value.toLowerCase()) + return ['y', 'yes', 't', 'true', 'on', '1'].includes(String(value).toLowerCase()) } export function determineNodeEnv(): NodeEnv { From 10eb8c492f72027a49fe56ed45f7ff79a8a93faa Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Thu, 29 Apr 2021 15:36:59 +0200 Subject: [PATCH 18/22] separate url for graphile retry queue, otherwise use existing postgres pool (fixes helm connection string issue) --- src/main/retry/graphile-queue.ts | 14 ++++++++++---- src/shared/config.ts | 2 ++ src/types.ts | 1 + 3 files changed, 13 insertions(+), 4 deletions(-) diff --git a/src/main/retry/graphile-queue.ts b/src/main/retry/graphile-queue.ts index e1cbbc53..50fe7add 100644 --- a/src/main/retry/graphile-queue.ts +++ b/src/main/retry/graphile-queue.ts @@ -1,4 +1,4 @@ -import { makeWorkerUtils, run, Runner, WorkerUtils } from 'graphile-worker' +import { makeWorkerUtils, run, Runner, WorkerUtils, WorkerUtilsOptions } from 'graphile-worker' import { EnqueuedRetry, OnRetryCallback, PluginsServer, RetryQueue } from '../../types' @@ -21,9 +21,15 @@ export class GraphileQueue implements RetryQueue { async enqueue(retry: EnqueuedRetry): Promise { if (!this.workerUtils) { - this.workerUtils = await makeWorkerUtils({ - connectionString: this.pluginsServer.DATABASE_URL, - }) + 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 }) diff --git a/src/shared/config.ts b/src/shared/config.ts index b2dba522..69b5446a 100644 --- a/src/shared/config.ts +++ b/src/shared/config.ts @@ -60,6 +60,7 @@ export function getDefaultConfig(): PluginsServerConfig { INTERNAL_MMDB_SERVER_PORT: 0, PLUGIN_SERVER_IDLE: false, RETRY_QUEUES: '', + RETRY_QUEUE_GRAPHILE_URL: '', ENABLE_PERSISTENT_CONSOLE: false, // TODO: remove when persistent console ships in main repo STALENESS_RESTART_SECONDS: 0, } @@ -104,6 +105,7 @@ export function getConfigHelp(): Record { 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', STALENESS_RESTART_SECONDS: 'trigger a restart if no event ingested for this duration', } } diff --git a/src/types.ts b/src/types.ts index 4ac1b272..5dd0f8b2 100644 --- a/src/types.ts +++ b/src/types.ts @@ -72,6 +72,7 @@ export interface PluginsServerConfig extends Record { INTERNAL_MMDB_SERVER_PORT: number PLUGIN_SERVER_IDLE: boolean RETRY_QUEUES: string + RETRY_QUEUE_GRAPHILE_URL: string ENABLE_PERSISTENT_CONSOLE: boolean STALENESS_RESTART_SECONDS: number } From a2e05b9c3ba93078a74015b87919b957d6c12dc0 Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Thu, 29 Apr 2021 17:50:40 +0200 Subject: [PATCH 19/22] convert startRedlock params to options object --- src/main/services/redlock.ts | 24 +++++++++++++++-------- src/main/services/retry-queue-consumer.ts | 12 ++++++------ src/main/services/schedule.ts | 12 ++++++------ 3 files changed, 28 insertions(+), 20 deletions(-) diff --git a/src/main/services/redlock.ts b/src/main/services/redlock.ts index fcf40df1..a28eac67 100644 --- a/src/main/services/redlock.ts +++ b/src/main/services/redlock.ts @@ -5,13 +5,21 @@ import { status } from '../../shared/status' import { createRedis } from '../../shared/utils' import { PluginsServer } from '../../types' -export async function startRedlock( - server: PluginsServer, - resource: string, - onLock: () => Promise | void, - onUnlock: () => Promise | void, - ttl = 60 -): Promise<() => Promise> { +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 @@ -19,7 +27,7 @@ export async function startRedlock( let lock: Redlock.Lock let lockTimeout: NodeJS.Timeout - const lockTTL = ttl * 1000 // 60 sec + const lockTTL = ttl * 1000 // 60 sec if default passed in const retryDelay = lockTTL / 10 // 6 sec const extendDelay = lockTTL / 2 // 30 sec diff --git a/src/main/services/retry-queue-consumer.ts b/src/main/services/retry-queue-consumer.ts index 11c710b9..b555bdff 100644 --- a/src/main/services/retry-queue-consumer.ts +++ b/src/main/services/retry-queue-consumer.ts @@ -20,19 +20,19 @@ export async function startRetryQueueConsumer( } } - const unlock = await startRedlock( + const unlock = await startRedlock({ server, - LOCKED_RESOURCE, - async () => { + resource: LOCKED_RESOURCE, + onLock: async () => { status.info('🔄', 'Retry queue consumer lock aquired') await server.retryQueueManager.startConsumer(onRetry) }, - async () => { + onUnlock: async () => { status.info('🔄', 'Stopping retry queue consumer') await server.retryQueueManager.stopConsumer() }, - server.SCHEDULE_LOCK_TTL - ) + 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 570abc98..2c34a550 100644 --- a/src/main/services/schedule.ts +++ b/src/main/services/schedule.ts @@ -40,18 +40,18 @@ export async function startSchedule( runTasksDebounced(server!, piscina!, 'runEveryDay') }) - const unlock = await startRedlock( + const unlock = await startRedlock({ server, - LOCKED_RESOURCE, - () => { + resource: LOCKED_RESOURCE, + onLock: () => { weHaveTheLock = true onLock?.() }, - () => { + onUnlock: () => { weHaveTheLock = false }, - server.SCHEDULE_LOCK_TTL - ) + ttl: server.SCHEDULE_LOCK_TTL, + }) const stopSchedule = async () => { stopped = true From a0266801410c8adc2dc016381846e6c1e3ffdc1e Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Thu, 29 Apr 2021 17:59:06 +0200 Subject: [PATCH 20/22] move type around --- src/main/retry/retry-queue-manager.ts | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/src/main/retry/retry-queue-manager.ts b/src/main/retry/retry-queue-manager.ts index c0b6d831..01514910 100644 --- a/src/main/retry/retry-queue-manager.ts +++ b/src/main/retry/retry-queue-manager.ts @@ -17,12 +17,12 @@ export class RetryQueueManager implements RetryQueue { this.pluginsServer = pluginsServer this.retryQueues = pluginsServer.RETRY_QUEUES.split(',') - .map((q) => q.trim()) + .map((q) => q.trim() as keyof typeof queues) .filter((q) => !!q) .map( (queue): RetryQueue => { - if (queues[queue as keyof typeof queues]) { - return queues[queue as keyof typeof queues](pluginsServer) + if (queues[queue]) { + return queues[queue](pluginsServer) } else { throw new Error(`Unknown retry queue "${queue}"`) } From 8d002d986a01f699363ede1a92a253361020fe0b Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Thu, 29 Apr 2021 18:02:54 +0200 Subject: [PATCH 21/22] use an enum --- src/main/retry/retry-queue-manager.ts | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/src/main/retry/retry-queue-manager.ts b/src/main/retry/retry-queue-manager.ts index 01514910..56c98cbb 100644 --- a/src/main/retry/retry-queue-manager.ts +++ b/src/main/retry/retry-queue-manager.ts @@ -4,7 +4,12 @@ import { EnqueuedRetry, OnRetryCallback, PluginsServer, RetryQueue } from '../.. import { FsQueue } from './fs-queue' import { GraphileQueue } from './graphile-queue' -const queues = { +enum QueueType { + FS = 'fs', + Graphile = 'graphile', +} + +const queues: Record RetryQueue> = { fs: () => new FsQueue(), graphile: (pluginsServer: PluginsServer) => new GraphileQueue(pluginsServer), } @@ -17,7 +22,7 @@ export class RetryQueueManager implements RetryQueue { this.pluginsServer = pluginsServer this.retryQueues = pluginsServer.RETRY_QUEUES.split(',') - .map((q) => q.trim() as keyof typeof queues) + .map((q) => q.trim() as QueueType) .filter((q) => !!q) .map( (queue): RetryQueue => { From 3990953059ed0565438aa3404a0f89ff468019ca Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Fri, 30 Apr 2021 14:50:49 +0200 Subject: [PATCH 22/22] update typo in comment --- src/main/services/redlock.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/main/services/redlock.ts b/src/main/services/redlock.ts index a28eac67..ce69dade 100644 --- a/src/main/services/redlock.ts +++ b/src/main/services/redlock.ts @@ -35,7 +35,7 @@ export async function startRedlock({ const redis = await createRedis(server) const redlock = new Redlock([redis], { - // we handle retires ourselves to have a way to cancel the promises on quit + // 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, })