diff --git a/benchmarks/clickhouse/e2e.kafka.benchmark.ts b/benchmarks/clickhouse/e2e.kafka.benchmark.ts index 0ab69288..07e3b1db 100644 --- a/benchmarks/clickhouse/e2e.kafka.benchmark.ts +++ b/benchmarks/clickhouse/e2e.kafka.benchmark.ts @@ -1,12 +1,12 @@ import { performance } from 'perf_hooks' -import { KAFKA_EVENTS_PLUGIN_INGESTION } from '../../src/ingestion/topics' -import { startPluginsServer } from '../../src/server' +import { startPluginsServer } from '../../src/main/pluginsServer' +import { KAFKA_EVENTS_PLUGIN_INGESTION } from '../../src/shared/ingestion/topics' +import { delay, UUIDT } from '../../src/shared/utils' import { LogLevel, PluginsServerConfig, Queue } from '../../src/types' import { PluginsServer } from '../../src/types' -import { delay, UUIDT } from '../../src/utils' -import { createPosthog, DummyPostHog } from '../../src/vm/extensions/posthog' import { makePiscina } from '../../src/worker/piscina' +import { createPosthog, DummyPostHog } from '../../src/worker/vm/extensions/posthog' import { resetTestDatabaseClickhouse } from '../../tests/helpers/clickhouse' import { resetKafka } from '../../tests/helpers/kafka' import { pluginConfig39 } from '../../tests/helpers/plugins' diff --git a/benchmarks/clickhouse/e2e.timeout.benchmark.ts b/benchmarks/clickhouse/e2e.timeout.benchmark.ts index ca18324e..f8c43ec7 100644 --- a/benchmarks/clickhouse/e2e.timeout.benchmark.ts +++ b/benchmarks/clickhouse/e2e.timeout.benchmark.ts @@ -1,12 +1,12 @@ import { performance } from 'perf_hooks' -import { KAFKA_EVENTS_PLUGIN_INGESTION } from '../../src/ingestion/topics' -import { startPluginsServer } from '../../src/server' +import { startPluginsServer } from '../../src/main/pluginsServer' +import { KAFKA_EVENTS_PLUGIN_INGESTION } from '../../src/shared/ingestion/topics' +import { delay, UUIDT } from '../../src/shared/utils' import { ClickHouseEvent, LogLevel, PluginsServerConfig, Queue } from '../../src/types' import { PluginsServer } from '../../src/types' -import { delay, UUIDT } from '../../src/utils' -import { createPosthog, DummyPostHog } from '../../src/vm/extensions/posthog' import { makePiscina } from '../../src/worker/piscina' +import { createPosthog, DummyPostHog } from '../../src/worker/vm/extensions/posthog' import { resetTestDatabaseClickhouse } from '../../tests/helpers/clickhouse' import { resetKafka } from '../../tests/helpers/kafka' import { pluginConfig39 } from '../../tests/helpers/plugins' diff --git a/benchmarks/postgres/e2e.celery.benchmark.ts b/benchmarks/postgres/e2e.celery.benchmark.ts index f647c80d..cf11be2a 100644 --- a/benchmarks/postgres/e2e.celery.benchmark.ts +++ b/benchmarks/postgres/e2e.celery.benchmark.ts @@ -1,12 +1,12 @@ import * as IORedis from 'ioredis' import { performance } from 'perf_hooks' -import { startPluginsServer } from '../../src/server' +import { startPluginsServer } from '../../src/main/pluginsServer' +import { delay, UUIDT } from '../../src/shared/utils' import { LogLevel, PluginsServerConfig, Queue } from '../../src/types' import { PluginsServer } from '../../src/types' -import { delay, UUIDT } from '../../src/utils' -import { createPosthog, DummyPostHog } from '../../src/vm/extensions/posthog' import { makePiscina } from '../../src/worker/piscina' +import { createPosthog, DummyPostHog } from '../../src/worker/vm/extensions/posthog' import { pluginConfig39 } from '../../tests/helpers/plugins' import { resetTestDatabase } from '../../tests/helpers/sql' import { delayUntilEventIngested } from '../../tests/shared/process-event' diff --git a/benchmarks/postgres/helpers/piscina.ts b/benchmarks/postgres/helpers/piscina.ts index 589b8c87..6ecc0143 100644 --- a/benchmarks/postgres/helpers/piscina.ts +++ b/benchmarks/postgres/helpers/piscina.ts @@ -1,9 +1,9 @@ import Piscina from '@posthog/piscina' import { PluginEvent } from '@posthog/plugin-scaffold/src/types' -import { defaultConfig } from '../../../src/config' +import { defaultConfig } from '../../../src/shared/config' +import { UUIDT } from '../../../src/shared/utils' import { LogLevel } from '../../../src/types' -import { UUIDT } from '../../../src/utils' import { makePiscina } from '../../../src/worker/piscina' export function setupPiscina(workers: number, tasksPerWorker: number): Piscina { diff --git a/benchmarks/postgres/ingestion.benchmark.ts b/benchmarks/postgres/ingestion.benchmark.ts index 8ead9203..6a58637c 100644 --- a/benchmarks/postgres/ingestion.benchmark.ts +++ b/benchmarks/postgres/ingestion.benchmark.ts @@ -4,15 +4,15 @@ import os from 'os' import { performance } from 'perf_hooks' import { IEvent } from '../../src/idl/protos' -import { EventsProcessor } from '../../src/ingestion/process-event' -import { createServer } from '../../src/server' +import { createServer } from '../../src/shared/server' +import { UUIDT } from '../../src/shared/utils' import { LogLevel, PluginsServer, SessionRecordingEvent, Team } from '../../src/types' -import { UUIDT } from '../../src/utils' +import { EventsProcessor } from '../../src/worker/ingestion/process-event' import { getFirstTeam, resetTestDatabase } from '../../tests/helpers/sql' import { endLog, startLog } from './helpers/log' import { ingestCountEvents, setupPiscina } from './helpers/piscina' -jest.mock('../../src/sql') +jest.mock('../../src/shared/sql') jest.setTimeout(600000) // 600 sec timeout describe('ingestion benchmarks', () => { diff --git a/benchmarks/vm/memory.benchmark.ts b/benchmarks/vm/memory.benchmark.ts index 412bfa2c..f2d1a422 100644 --- a/benchmarks/vm/memory.benchmark.ts +++ b/benchmarks/vm/memory.benchmark.ts @@ -1,11 +1,11 @@ import { PluginEvent } from '@posthog/plugin-scaffold/src/types' -import { createServer } from '../../src/server' +import { createServer } from '../../src/shared/server' import { Plugin, PluginConfig, PluginConfigVMReponse } from '../../src/types' -import { createPluginConfigVM } from '../../src/vm/vm' +import { createPluginConfigVM } from '../../src/worker/vm/vm' import { commonOrganizationId } from '../../tests/helpers/plugins' -jest.mock('../../src/sql') +jest.mock('../../src/shared/sql') jest.setTimeout(600000) // 600 sec timeout function createEvent(index: number): PluginEvent { diff --git a/benchmarks/vm/worker.benchmark.ts b/benchmarks/vm/worker.benchmark.ts index 4c75e504..ed824b14 100644 --- a/benchmarks/vm/worker.benchmark.ts +++ b/benchmarks/vm/worker.benchmark.ts @@ -2,12 +2,12 @@ import { PluginEvent } from '@posthog/plugin-scaffold/src/types' import * as os from 'os' import { performance } from 'perf_hooks' -import { defaultConfig } from '../../src/config' +import { defaultConfig } from '../../src/shared/config' import { LogLevel } from '../../src/types' import { makePiscina } from '../../src/worker/piscina' import { resetTestDatabase } from '../../tests/helpers/sql' -jest.mock('../../src/sql') +jest.mock('../../src/shared/sql') jest.setTimeout(600000) // 600 sec timeout function processOneEvent( diff --git a/src/healthcheck.ts b/src/healthcheck.ts index 9a684a69..90e96a19 100644 --- a/src/healthcheck.ts +++ b/src/healthcheck.ts @@ -1,7 +1,7 @@ import Redis from 'ioredis' -import { defaultConfig } from './config' -import { Status } from './status' +import { defaultConfig } from './shared/config' +import { Status } from './shared/status' const healthStatus = new Status('HLTH') diff --git a/src/index.ts b/src/index.ts index 71a6ec5e..0cb9164f 100644 --- a/src/index.ts +++ b/src/index.ts @@ -1,8 +1,8 @@ import * as yargs from 'yargs' -import { configHelp, defaultConfig } from './config' import { initApp } from './init' -import { startPluginsServer } from './server' +import { startPluginsServer } from './main/pluginsServer' +import { configHelp, defaultConfig } from './shared/config' import { PluginsServerConfig } from './types' import { makePiscina } from './worker/piscina' diff --git a/src/init.ts b/src/init.ts index d7e1948d..6ff02452 100644 --- a/src/init.ts +++ b/src/init.ts @@ -1,7 +1,7 @@ import * as Sentry from '@sentry/node' +import { setLogLevel } from './shared/utils' import { PluginsServerConfig } from './types' -import { setLogLevel } from './utils' // Must require as `tsc` strips unused `import` statements and just requiring this seems to init some globals require('@sentry/tracing') diff --git a/src/celery/worker.ts b/src/main/ingestion/celery-queue-worker.ts similarity index 96% rename from src/celery/worker.ts rename to src/main/ingestion/celery-queue-worker.ts index 0fa7c1dc..95c2a57b 100644 --- a/src/celery/worker.ts +++ b/src/main/ingestion/celery-queue-worker.ts @@ -1,11 +1,11 @@ -import { status } from '../status' -import { Queue } from '../types' -import Base from './base' -import { Message } from './message' +import Base from '../../shared/celery/base' +import { Message } from '../../shared/celery/message' +import { status } from '../../shared/status' +import { Queue } from '../../types' type Handler = (...args: any[]) => Promise -export class Worker extends Base implements Queue { +export class CeleryQueueWorker extends Base implements Queue { handlers: Record = {} activeTasks: Set> = new Set() @@ -225,4 +225,4 @@ export class Worker extends Base implements Queue { } } -export default Worker +export default CeleryQueueWorker diff --git a/src/ingestion/kafka-queue.ts b/src/main/ingestion/kafka-queue.ts similarity index 98% rename from src/ingestion/kafka-queue.ts rename to src/main/ingestion/kafka-queue.ts index 469b55d8..13b96712 100644 --- a/src/ingestion/kafka-queue.ts +++ b/src/main/ingestion/kafka-queue.ts @@ -3,9 +3,9 @@ import * as Sentry from '@sentry/node' import { Consumer, EachBatchPayload, Kafka } from 'kafkajs' import { PluginsServer, Queue } from 'types' -import { status } from '../status' -import { groupIntoBatches, killGracefully } from '../utils' -import { timeoutGuard } from './utils' +import { timeoutGuard } from '../../shared/ingestion/utils' +import { status } from '../../shared/status' +import { groupIntoBatches, killGracefully } from '../../shared/utils' export class KafkaQueue implements Queue { private pluginsServer: PluginsServer diff --git a/src/main/pluginsServer.ts b/src/main/pluginsServer.ts new file mode 100644 index 00000000..721a0a21 --- /dev/null +++ b/src/main/pluginsServer.ts @@ -0,0 +1,149 @@ +import Piscina from '@posthog/piscina' +import * as Sentry from '@sentry/node' +import { FastifyInstance } from 'fastify' +import Redis from 'ioredis' +import * as schedule from 'node-schedule' + +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 { startQueue } from './queue' +import { startSchedule } from './services/schedule' +import { startFastifyInstance, stopFastifyInstance } from './web/server' + +const { version } = require('../../package.json') + +// TODO: refactor this into a class, removing the need for many different Servers +export type ServerInstance = { + server: PluginsServer + piscina: Piscina + queue: Queue + stop: () => Promise +} + +export async function startPluginsServer( + config: Partial, + makePiscina: (config: PluginsServerConfig) => Piscina +): Promise { + const serverConfig: PluginsServerConfig = { + ...defaultConfig, + ...config, + } + + status.info('⚡', `posthog-plugin-server v${version}`) + status.info('â„šī¸', `${serverConfig.WORKER_CONCURRENCY} workers, ${serverConfig.TASKS_PER_WORKER} tasks per worker`) + + let pubSub: Redis.Redis | undefined + let server: PluginsServer | undefined + let fastifyInstance: FastifyInstance | undefined + let pingJob: schedule.Job | undefined + let statsJob: schedule.Job | undefined + let piscina: Piscina | undefined + let queue: Queue | undefined + let closeServer: () => Promise | undefined + let scheduleControl: ScheduleControl | undefined + + let shutdownStatus = 0 + + async function closeJobs(): Promise { + shutdownStatus += 1 + if (shutdownStatus === 2) { + status.info('🔁', 'Try again to shut down forcibly') + return + } + if (shutdownStatus >= 3) { + status.info('â—ī¸', 'Shutting down forcibly!') + void piscina?.destroy() + process.exit() + } + status.info('💤', ' Shutting down gracefully...') + if (fastifyInstance && !serverConfig?.DISABLE_WEB) { + await stopFastifyInstance(fastifyInstance!) + } + await queue?.stop() + await pubSub?.quit() + pingJob && schedule.cancelJob(pingJob) + statsJob && schedule.cancelJob(statsJob) + await scheduleControl?.stopSchedule() + if (piscina) { + await stopPiscina(piscina) + } + await closeServer?.() + + // wait an extra second for any misc async task to finish + await delay(1000) + } + + for (const signal of ['SIGINT', 'SIGTERM', 'SIGHUP']) { + process.on(signal, closeJobs) + } + + try { + ;[server, closeServer] = await createServer(serverConfig, null) + + piscina = makePiscina(serverConfig) + if (!server.DISABLE_WEB) { + fastifyInstance = await startFastifyInstance(server) + } + + scheduleControl = await startSchedule(server, piscina) + queue = await startQueue(server, piscina) + piscina.on('drain', () => { + queue?.resume() + }) + + // use one extra connection for redis pubsub + pubSub = await createRedis(server) + await pubSub.subscribe(server.PLUGINS_RELOAD_PUBSUB_CHANNEL) + pubSub.on('message', async (channel: string, message) => { + if (channel === server!.PLUGINS_RELOAD_PUBSUB_CHANNEL) { + status.info('⚡', 'Reloading plugins!') + + await piscina?.broadcastTask({ task: 'reloadPlugins' }) + await scheduleControl?.reloadSchedule() + } + }) + + // every 5 seconds set Redis keys @posthog-plugin-server/ping and @posthog-plugin-server/version + pingJob = schedule.scheduleJob('*/5 * * * * *', async () => { + await server!.db!.redisSet('@posthog-plugin-server/ping', new Date().toISOString(), 60, { + jsonSerialize: false, + }) + await server!.db!.redisSet('@posthog-plugin-server/version', version, undefined, { jsonSerialize: false }) + }) + + // every 10 seconds sends stuff to StatsD + statsJob = schedule.scheduleJob('*/10 * * * * *', () => { + if (piscina) { + server!.statsd?.gauge(`piscina.utilization`, (piscina?.utilization || 0) * 100) + server!.statsd?.gauge(`piscina.threads`, piscina?.threads.length) + server!.statsd?.gauge(`piscina.queue_size`, piscina?.queueSize) + } + }) + + status.info('🚀', 'All systems go.') + } catch (error) { + Sentry.captureException(error) + status.error('đŸ’Ĩ', 'Launchpad failure!', error) + void Sentry.flush() // flush in the background + await closeJobs() + process.exit(1) + } + + return { + server, + piscina, + queue, + stop: closeJobs, + } +} + +export async function stopPiscina(piscina: Piscina): Promise { + // Wait two seconds for any running workers to stop. + // TODO: better "wait until everything is done" + await delay(2000) + await Promise.race([piscina.broadcastTask({ task: 'flushKafkaMessages' }), delay(2000)]) + await piscina.destroy() +} diff --git a/src/worker/queue.ts b/src/main/queue.ts similarity index 89% rename from src/worker/queue.ts rename to src/main/queue.ts index 97e3ccb3..5e6fc70d 100644 --- a/src/worker/queue.ts +++ b/src/main/queue.ts @@ -2,13 +2,12 @@ import Piscina from '@posthog/piscina' import { PluginEvent } from '@posthog/plugin-scaffold' import * as Sentry from '@sentry/node' -import Client from '../celery/client' -import Worker from '../celery/worker' -import { IngestEventResponse } from '../ingestion/ingest-event' -import { KafkaQueue } from '../ingestion/kafka-queue' -import { status } from '../status' -import { PluginsServer, Queue } from '../types' -import { UUIDT } from '../utils' +import Client from '../shared/celery/client' +import { status } from '../shared/status' +import { UUIDT } from '../shared/utils' +import { IngestEventResponse, PluginsServer, Queue } from '../types' +import CeleryQueueWorker from './ingestion/celery-queue-worker' +import { KafkaQueue } from './ingestion/kafka-queue' export type WorkerMethods = { processEvent: (event: PluginEvent) => Promise @@ -53,7 +52,7 @@ export async function startQueue( } function startQueueRedis(server: PluginsServer, piscina: Piscina | undefined, workerMethods: WorkerMethods): Queue { - const celeryQueue = new Worker(server.db, server.PLUGINS_CELERY_QUEUE) + const celeryQueue = new CeleryQueueWorker(server.db, server.PLUGINS_CELERY_QUEUE) const client = new Client(server.db, server.CELERY_DEFAULT_QUEUE) celeryQueue.register( diff --git a/src/services/schedule.ts b/src/main/services/schedule.ts similarity index 97% rename from src/services/schedule.ts rename to src/main/services/schedule.ts index 94c6014b..a1ac9ce3 100644 --- a/src/services/schedule.ts +++ b/src/main/services/schedule.ts @@ -3,10 +3,10 @@ import * as Sentry from '@sentry/node' import * as schedule from 'node-schedule' import Redlock from 'redlock' -import { processError } from '../error' -import { status } from '../status' -import { PluginConfigId, PluginsServer, ScheduleControl } from '../types' -import { createRedis, delay } from '../utils' +import { processError } from '../../shared/error' +import { status } from '../../shared/status' +import { createRedis, delay } from '../../shared/utils' +import { PluginConfigId, PluginsServer, ScheduleControl } from '../../types' export const LOCKED_RESOURCE = 'plugin-server:locks:schedule' diff --git a/src/web/server.ts b/src/main/web/server.ts similarity index 95% rename from src/web/server.ts rename to src/main/web/server.ts index 7effae00..ffdde008 100644 --- a/src/web/server.ts +++ b/src/main/web/server.ts @@ -1,7 +1,7 @@ import { fastify, FastifyInstance } from 'fastify' import { PluginsServer } from 'types' -import { status } from '../status' +import { status } from '../../shared/status' export function buildFastifyInstance(): FastifyInstance { const fastifyInstance = fastify() diff --git a/src/celery/base.ts b/src/shared/celery/base.ts similarity index 100% rename from src/celery/base.ts rename to src/shared/celery/base.ts diff --git a/src/celery/broker.ts b/src/shared/celery/broker.ts similarity index 99% rename from src/celery/broker.ts rename to src/shared/celery/broker.ts index 83308fb6..1ae60b1c 100644 --- a/src/celery/broker.ts +++ b/src/shared/celery/broker.ts @@ -1,8 +1,8 @@ import { v4 } from 'uuid' +import { Pausable } from '../../types' import { DB } from '../db' import { status } from '../status' -import { Pausable } from '../types' import { Message } from './message' type BrokerSubscription = { queue: string; callback: (message: Message) => any } diff --git a/src/celery/client.ts b/src/shared/celery/client.ts similarity index 100% rename from src/celery/client.ts rename to src/shared/celery/client.ts diff --git a/src/celery/conf.ts b/src/shared/celery/conf.ts similarity index 100% rename from src/celery/conf.ts rename to src/shared/celery/conf.ts diff --git a/src/celery/message.ts b/src/shared/celery/message.ts similarity index 100% rename from src/celery/message.ts rename to src/shared/celery/message.ts diff --git a/src/celery/task.ts b/src/shared/celery/task.ts similarity index 100% rename from src/celery/task.ts rename to src/shared/celery/task.ts diff --git a/src/config.ts b/src/shared/config.ts similarity index 98% rename from src/config.ts rename to src/shared/config.ts index 0b1a0588..b3671c73 100644 --- a/src/config.ts +++ b/src/shared/config.ts @@ -1,7 +1,7 @@ import os from 'os' +import { LogLevel, PluginsServerConfig } from '../types' import { KAFKA_EVENTS_PLUGIN_INGESTION } from './ingestion/topics' -import { LogLevel, PluginsServerConfig } from './types' export const defaultConfig = overrideWithEnv(getDefaultConfig()) export const configHelp = getConfigHelp() diff --git a/src/db.ts b/src/shared/db.ts similarity index 99% rename from src/db.ts rename to src/shared/db.ts index 34bff1af..6c2f17e8 100644 --- a/src/db.ts +++ b/src/shared/db.ts @@ -7,8 +7,6 @@ import { Producer, ProducerRecord } from 'kafkajs' import { DateTime } from 'luxon' import { Pool, PoolClient, QueryConfig, QueryResult, QueryResultRow } from 'pg' -import { KAFKA_PERSON, KAFKA_PERSON_UNIQUE_ID } from './ingestion/topics' -import { chainToElements, hashElements, timeoutGuard, unparsePersonPartial } from './ingestion/utils' import { ClickHouseEvent, ClickHousePerson, @@ -25,7 +23,9 @@ import { RawPerson, SessionRecordingEvent, TimestampFormat, -} from './types' +} from '../types' +import { KAFKA_PERSON, KAFKA_PERSON_UNIQUE_ID } from './ingestion/topics' +import { chainToElements, hashElements, timeoutGuard, unparsePersonPartial } from './ingestion/utils' import { castTimestampOrNow, clickHouseTimestampToISO, diff --git a/src/error.ts b/src/shared/error.ts similarity index 98% rename from src/error.ts rename to src/shared/error.ts index 7a57a763..d8a83ad0 100644 --- a/src/error.ts +++ b/src/shared/error.ts @@ -1,7 +1,7 @@ import { PluginEvent } from '@posthog/plugin-scaffold' +import { PluginConfig, PluginConfigId, PluginError, PluginsServer } from '../types' import { setError } from './sql' -import { PluginConfig, PluginConfigId, PluginError, PluginsServer } from './types' export async function processError( server: PluginsServer, diff --git a/src/ingestion/topics.ts b/src/shared/ingestion/topics.ts similarity index 100% rename from src/ingestion/topics.ts rename to src/shared/ingestion/topics.ts diff --git a/src/ingestion/utils.ts b/src/shared/ingestion/utils.ts similarity index 99% rename from src/ingestion/utils.ts rename to src/shared/ingestion/utils.ts index fba7c96d..8d0bf47d 100644 --- a/src/ingestion/utils.ts +++ b/src/shared/ingestion/utils.ts @@ -2,8 +2,8 @@ import { Properties } from '@posthog/plugin-scaffold' import * as Sentry from '@sentry/node' import crypto from 'crypto' +import { BasePerson, Element, Person, RawPerson } from '../../types' import { defaultConfig } from '../config' -import { BasePerson, Element, Person, RawPerson } from '../types' export function unparsePersonPartial(person: Partial): Partial { return { ...(person as BasePerson), ...(person.created_at ? { created_at: person.created_at.toISO() } : {}) } diff --git a/src/server.ts b/src/shared/server.ts similarity index 53% rename from src/server.ts rename to src/shared/server.ts index 9bf727b1..afa9f58f 100644 --- a/src/server.ts +++ b/src/shared/server.ts @@ -1,30 +1,23 @@ import ClickHouse from '@posthog/clickhouse' -import Piscina from '@posthog/piscina' -import { PluginEvent } from '@posthog/plugin-scaffold' import * as Sentry from '@sentry/node' -import { FastifyInstance } from 'fastify' import * as fs from 'fs' import { createPool } from 'generic-pool' import { StatsD } from 'hot-shots' import Redis from 'ioredis' import { Kafka, logLevel, Producer } from 'kafkajs' import { DateTime } from 'luxon' -import * as schedule from 'node-schedule' import * as path from 'path' -import { Pool, types as pgTypes } from 'pg' +import { types as pgTypes } from 'pg' import { ConnectionOptions } from 'tls' +import { PluginsServer, PluginsServerConfig } from '../types' +import { EventsProcessor } from '../worker/ingestion/process-event' import { defaultConfig } from './config' import { DB } from './db' -import { EventsProcessor } from './ingestion/process-event' -import { startSchedule } from './services/schedule' import { status } from './status' -import { PluginsServer, PluginsServerConfig, Queue, ScheduleControl } from './types' -import { createPostgresPool, createRedis, delay, UUIDT } from './utils' -import { startFastifyInstance, stopFastifyInstance } from './web/server' -import { startQueue } from './worker/queue' +import { createPostgresPool, createRedis, UUIDT } from './utils' -const { version } = require('../package.json') +const { version } = require('../../package.json') export async function createServer( config: Partial = {}, @@ -160,6 +153,7 @@ export async function createServer( pluginSchedulePromises: { runEveryMinute: {}, runEveryHour: {}, runEveryDay: {} }, } + // :TODO: This is only used on worker threads, not main server.eventsProcessor = new EventsProcessor(server as PluginsServer) const closeServer = async () => { @@ -173,136 +167,3 @@ export async function createServer( return [server as PluginsServer, closeServer] } - -// TODO: refactor this into a class, removing the need for many different Servers -export type ServerInstance = { - server: PluginsServer - piscina: Piscina - queue: Queue - stop: () => Promise -} - -export async function startPluginsServer( - config: Partial, - makePiscina: (config: PluginsServerConfig) => Piscina -): Promise { - const serverConfig: PluginsServerConfig = { - ...defaultConfig, - ...config, - } - - status.info('⚡', `posthog-plugin-server v${version}`) - status.info('â„šī¸', `${serverConfig.WORKER_CONCURRENCY} workers, ${serverConfig.TASKS_PER_WORKER} tasks per worker`) - - let pubSub: Redis.Redis | undefined - let server: PluginsServer | undefined - let fastifyInstance: FastifyInstance | undefined - let pingJob: schedule.Job | undefined - let statsJob: schedule.Job | undefined - let piscina: Piscina | undefined - let queue: Queue | undefined - let closeServer: () => Promise | undefined - let scheduleControl: ScheduleControl | undefined - - let shutdownStatus = 0 - - async function closeJobs(): Promise { - shutdownStatus += 1 - if (shutdownStatus === 2) { - status.info('🔁', 'Try again to shut down forcibly') - return - } - if (shutdownStatus >= 3) { - status.info('â—ī¸', 'Shutting down forcibly!') - void piscina?.destroy() - process.exit() - } - status.info('💤', ' Shutting down gracefully...') - if (fastifyInstance && !serverConfig?.DISABLE_WEB) { - await stopFastifyInstance(fastifyInstance!) - } - await queue?.stop() - await pubSub?.quit() - pingJob && schedule.cancelJob(pingJob) - statsJob && schedule.cancelJob(statsJob) - await scheduleControl?.stopSchedule() - if (piscina) { - await stopPiscina(piscina) - } - await closeServer?.() - - // wait an extra second for any misc async task to finish - await delay(1000) - } - - for (const signal of ['SIGINT', 'SIGTERM', 'SIGHUP']) { - process.on(signal, closeJobs) - } - - try { - ;[server, closeServer] = await createServer(serverConfig, null) - - piscina = makePiscina(serverConfig) - if (!server.DISABLE_WEB) { - fastifyInstance = await startFastifyInstance(server) - } - - scheduleControl = await startSchedule(server, piscina) - queue = await startQueue(server, piscina) - piscina.on('drain', () => { - queue?.resume() - }) - - // use one extra connection for redis pubsub - pubSub = await createRedis(server) - await pubSub.subscribe(server.PLUGINS_RELOAD_PUBSUB_CHANNEL) - pubSub.on('message', async (channel: string, message) => { - if (channel === server!.PLUGINS_RELOAD_PUBSUB_CHANNEL) { - status.info('⚡', 'Reloading plugins!') - - await piscina?.broadcastTask({ task: 'reloadPlugins' }) - await scheduleControl?.reloadSchedule() - } - }) - - // every 5 seconds set Redis keys @posthog-plugin-server/ping and @posthog-plugin-server/version - pingJob = schedule.scheduleJob('*/5 * * * * *', async () => { - await server!.db!.redisSet('@posthog-plugin-server/ping', new Date().toISOString(), 60, { - jsonSerialize: false, - }) - await server!.db!.redisSet('@posthog-plugin-server/version', version, undefined, { jsonSerialize: false }) - }) - - // every 10 seconds sends stuff to StatsD - statsJob = schedule.scheduleJob('*/10 * * * * *', () => { - if (piscina) { - server!.statsd?.gauge(`piscina.utilization`, (piscina?.utilization || 0) * 100) - server!.statsd?.gauge(`piscina.threads`, piscina?.threads.length) - server!.statsd?.gauge(`piscina.queue_size`, piscina?.queueSize) - } - }) - - status.info('🚀', 'All systems go.') - } catch (error) { - Sentry.captureException(error) - status.error('đŸ’Ĩ', 'Launchpad failure!', error) - void Sentry.flush() // flush in the background - await closeJobs() - process.exit(1) - } - - return { - server, - piscina, - queue, - stop: closeJobs, - } -} - -export async function stopPiscina(piscina: Piscina): Promise { - // Wait two seconds for any running workers to stop. - // TODO: better "wait until everything is done" - await delay(2000) - await Promise.race([piscina.broadcastTask({ task: 'flushKafkaMessages' }), delay(2000)]) - await piscina.destroy() -} diff --git a/src/sql.ts b/src/shared/sql.ts similarity index 97% rename from src/sql.ts rename to src/shared/sql.ts index f6246dc4..a4ff0c88 100644 --- a/src/sql.ts +++ b/src/shared/sql.ts @@ -1,4 +1,4 @@ -import { Plugin, PluginAttachmentDB, PluginConfig, PluginConfigId, PluginError, PluginsServer } from './types' +import { Plugin, PluginAttachmentDB, PluginConfig, PluginConfigId, PluginError, PluginsServer } from '../types' function pluginConfigsInForceQuery(specificField?: keyof PluginConfig): string { return `SELECT posthog_pluginconfig.${specificField || '*'} diff --git a/src/status.ts b/src/shared/status.ts similarity index 100% rename from src/status.ts rename to src/shared/status.ts diff --git a/src/utils.ts b/src/shared/utils.ts similarity index 99% rename from src/utils.ts rename to src/shared/utils.ts index 2b241d15..8aca9455 100644 --- a/src/utils.ts +++ b/src/shared/utils.ts @@ -8,8 +8,8 @@ import { Readable } from 'stream' import * as tar from 'tar-stream' import * as zlib from 'zlib' +import { LogLevel, Plugin, PluginsServerConfig, TimestampFormat } from '../types' import { status } from './status' -import { LogLevel, Plugin, PluginsServerConfig, TimestampFormat } from './types' /** Time until autoexit (due to error) gives up on graceful exit and kills the process right away. */ const GRACEFUL_EXIT_PERIOD_SECONDS = 5 diff --git a/src/types.ts b/src/types.ts index e2bf5629..5cf8f929 100644 --- a/src/types.ts +++ b/src/types.ts @@ -2,15 +2,15 @@ import ClickHouse from '@posthog/clickhouse' import { PluginAttachment, PluginConfigSchema, PluginEvent, Properties } from '@posthog/plugin-scaffold' import { Pool as GenericPool } from 'generic-pool' import { StatsD } from 'hot-shots' -import { EventsProcessor } from 'ingestion/process-event' import { Redis } from 'ioredis' import { Kafka, Producer } from 'kafkajs' import { DateTime } from 'luxon' import { Pool } from 'pg' import { VM } from 'vm2' -import { DB } from './db' -import { LazyPluginVM } from './vm/lazy' +import { DB } from './shared/db' +import { EventsProcessor } from './worker/ingestion/process-event' +import { LazyPluginVM } from './worker/vm/lazy' export enum LogLevel { Debug = 'debug', @@ -361,3 +361,5 @@ export interface ScheduleControl { stopSchedule: () => Promise reloadSchedule: () => Promise } + +export type IngestEventResponse = { success?: boolean; error?: string } diff --git a/src/ingestion/ingest-event.ts b/src/worker/ingestion/ingest-event.ts similarity index 86% rename from src/ingestion/ingest-event.ts rename to src/worker/ingestion/ingest-event.ts index 46e83158..6bee1f82 100644 --- a/src/ingestion/ingest-event.ts +++ b/src/worker/ingestion/ingest-event.ts @@ -2,11 +2,9 @@ import { PluginEvent } from '@posthog/plugin-scaffold' import * as Sentry from '@sentry/node' import { DateTime } from 'luxon' -import { status } from '../status' -import { PluginsServer } from '../types' -import { timeoutGuard } from './utils' - -export type IngestEventResponse = { success?: boolean; error?: string } +import { timeoutGuard } from '../../shared/ingestion/utils' +import { status } from '../../shared/status' +import { IngestEventResponse, PluginsServer } from '../../types' export async function ingestEvent(server: PluginsServer, event: PluginEvent): Promise { const timeout = timeoutGuard('Still ingesting event inside worker. Timeout warning after 30 sec!', { diff --git a/src/ingestion/process-event.ts b/src/worker/ingestion/process-event.ts similarity index 97% rename from src/ingestion/process-event.ts rename to src/worker/ingestion/process-event.ts index 704facaf..01e18be1 100644 --- a/src/ingestion/process-event.ts +++ b/src/worker/ingestion/process-event.ts @@ -7,10 +7,18 @@ import { DateTime, Duration } from 'luxon' import * as fetch from 'node-fetch' import { nodePostHog } from 'posthog-js-lite/dist/src/targets/node' -import Client from '../celery/client' -import { DB } from '../db' -import { Event as EventProto, IEvent } from '../idl/protos' -import { status } from '../status' +import { Event as EventProto, IEvent } from '../../idl/protos' +import Client from '../../shared/celery/client' +import { DB } from '../../shared/db' +import { KAFKA_EVENTS, KAFKA_SESSION_RECORDING_EVENTS } from '../../shared/ingestion/topics' +import { + elementsToString, + personInitialAndUTMProperties, + sanitizeEventName, + timeoutGuard, +} from '../../shared/ingestion/utils' +import { status } from '../../shared/status' +import { castTimestampOrNow, UUID, UUIDT } from '../../shared/utils' import { CohortPeople, Element, @@ -21,10 +29,7 @@ import { SessionRecordingEvent, Team, TimestampFormat, -} from '../types' -import { castTimestampOrNow, UUID, UUIDT } from '../utils' -import { KAFKA_EVENTS, KAFKA_SESSION_RECORDING_EVENTS } from './topics' -import { elementsToString, personInitialAndUTMProperties, sanitizeEventName, timeoutGuard } from './utils' +} from '../../types' export class EventsProcessor { pluginsServer: PluginsServer diff --git a/src/plugins/loadPlugin.ts b/src/worker/plugins/loadPlugin.ts similarity index 96% rename from src/plugins/loadPlugin.ts rename to src/worker/plugins/loadPlugin.ts index e0bc2c74..7cbe9ac8 100644 --- a/src/plugins/loadPlugin.ts +++ b/src/worker/plugins/loadPlugin.ts @@ -1,9 +1,9 @@ import * as fs from 'fs' import * as path from 'path' -import { processError } from '../error' -import { PluginConfig, PluginJsonConfig, PluginsServer } from '../types' -import { getFileFromArchive, pluginDigest } from '../utils' +import { processError } from '../../shared/error' +import { getFileFromArchive, pluginDigest } from '../../shared/utils' +import { PluginConfig, PluginJsonConfig, PluginsServer } from '../../types' export async function loadPlugin(server: PluginsServer, pluginConfig: PluginConfig): Promise { const { plugin } = pluginConfig diff --git a/src/plugins/run.ts b/src/worker/plugins/run.ts similarity index 96% rename from src/plugins/run.ts rename to src/worker/plugins/run.ts index 897b5874..78b3bb59 100644 --- a/src/plugins/run.ts +++ b/src/worker/plugins/run.ts @@ -1,7 +1,7 @@ import { PluginEvent } from '@posthog/plugin-scaffold' -import { processError } from '../error' -import { PluginConfig, PluginsServer } from '../types' +import { processError } from '../../shared/error' +import { PluginConfig, PluginsServer } from '../../types' export async function runPlugins(server: PluginsServer, event: PluginEvent): Promise { const pluginsToRun = getPluginsForTeam(server, event.team_id) diff --git a/src/plugins/setup.ts b/src/worker/plugins/setup.ts similarity index 97% rename from src/plugins/setup.ts rename to src/worker/plugins/setup.ts index 777bdd20..8e95cd85 100644 --- a/src/plugins/setup.ts +++ b/src/worker/plugins/setup.ts @@ -1,8 +1,8 @@ import { PluginAttachment } from '@posthog/plugin-scaffold' -import { getPluginAttachmentRows, getPluginConfigRows, getPluginRows } from '../sql' -import { status } from '../status' -import { Plugin, PluginConfig, PluginConfigId, PluginId, PluginsServer, TeamId } from '../types' +import { getPluginAttachmentRows, getPluginConfigRows, getPluginRows } from '../../shared/sql' +import { status } from '../../shared/status' +import { Plugin, PluginConfig, PluginConfigId, PluginId, PluginsServer, TeamId } from '../../types' import { LazyPluginVM } from '../vm/lazy' import { loadPlugin } from './loadPlugin' diff --git a/src/vm/extensions/cache.ts b/src/worker/vm/extensions/cache.ts similarity index 94% rename from src/vm/extensions/cache.ts rename to src/worker/vm/extensions/cache.ts index 00911a6c..ca19b5d2 100644 --- a/src/vm/extensions/cache.ts +++ b/src/worker/vm/extensions/cache.ts @@ -1,6 +1,6 @@ import { CacheExtension } from '@posthog/plugin-scaffold' -import { PluginsServer } from '../../types' +import { PluginsServer } from '../../../types' export function createCache(server: PluginsServer, pluginId: number, teamId: number): CacheExtension { const getKey = (key: string) => `@plugin/${pluginId}/${typeof teamId === 'undefined' ? '@all' : teamId}/${key}` diff --git a/src/vm/extensions/console.ts b/src/worker/vm/extensions/console.ts similarity index 100% rename from src/vm/extensions/console.ts rename to src/worker/vm/extensions/console.ts diff --git a/src/vm/extensions/google.ts b/src/worker/vm/extensions/google.ts similarity index 100% rename from src/vm/extensions/google.ts rename to src/worker/vm/extensions/google.ts diff --git a/src/vm/extensions/posthog.ts b/src/worker/vm/extensions/posthog.ts similarity index 94% rename from src/vm/extensions/posthog.ts rename to src/worker/vm/extensions/posthog.ts index 203d563d..4c057df0 100644 --- a/src/vm/extensions/posthog.ts +++ b/src/worker/vm/extensions/posthog.ts @@ -2,10 +2,10 @@ import { Properties } from '@posthog/plugin-scaffold' import { DateTime } from 'luxon' import { PluginConfig, PluginsServer, RawEventMessage } from 'types' -import Client from '../../celery/client' -import { UUIDT } from '../../utils' +import Client from '../../../shared/celery/client' +import { UUIDT } from '../../../shared/utils' -const { version } = require('../../../package.json') +const { version } = require('../../../../package.json') interface InternalData { distinct_id: string diff --git a/src/vm/extensions/storage.ts b/src/worker/vm/extensions/storage.ts similarity index 95% rename from src/vm/extensions/storage.ts rename to src/worker/vm/extensions/storage.ts index 80f862ac..a0b2f1a7 100644 --- a/src/vm/extensions/storage.ts +++ b/src/worker/vm/extensions/storage.ts @@ -1,6 +1,6 @@ import { StorageExtension } from '@posthog/plugin-scaffold' -import { PluginConfig, PluginsServer } from '../../types' +import { PluginConfig, PluginsServer } from '../../../types' export function createStorage(server: PluginsServer, pluginConfig: PluginConfig): StorageExtension { const get = async function (key: string, defaultValue: unknown): Promise { diff --git a/src/vm/lazy.ts b/src/worker/vm/lazy.ts similarity index 93% rename from src/vm/lazy.ts rename to src/worker/vm/lazy.ts index a0d81dd5..0a77f548 100644 --- a/src/vm/lazy.ts +++ b/src/worker/vm/lazy.ts @@ -1,6 +1,6 @@ -import { clearError, processError } from '../error' -import { status } from '../status' -import { PluginConfig, PluginConfigVMReponse, PluginsServer, PluginTask } from '../types' +import { clearError, processError } from '../../shared/error' +import { status } from '../../shared/status' +import { PluginConfig, PluginConfigVMReponse, PluginsServer, PluginTask } from '../../types' import { createPluginConfigVM } from './vm' export class LazyPluginVM { diff --git a/src/vm/transforms/common.ts b/src/worker/vm/transforms/common.ts similarity index 79% rename from src/vm/transforms/common.ts rename to src/worker/vm/transforms/common.ts index 6c60ed2a..32f1e785 100644 --- a/src/vm/transforms/common.ts +++ b/src/worker/vm/transforms/common.ts @@ -1,6 +1,6 @@ import { PluginObj } from '@babel/core' import * as types from '@babel/types' -import { PluginsServer } from '../../types' +import { PluginsServer } from '../../../types' export type PluginGen = (server: PluginsServer) => (param: { types: typeof types }) => PluginObj diff --git a/src/vm/transforms/index.ts b/src/worker/vm/transforms/index.ts similarity index 95% rename from src/vm/transforms/index.ts rename to src/worker/vm/transforms/index.ts index 10c62059..05ebe8c2 100644 --- a/src/vm/transforms/index.ts +++ b/src/worker/vm/transforms/index.ts @@ -1,6 +1,6 @@ import { transform } from '@babel/standalone' -import { PluginsServer } from '../../types' +import { PluginsServer } from '../../../types' import { loopTimeout } from './loop-timeout' import { promiseTimeout } from './promise-timeout' diff --git a/src/vm/transforms/loop-timeout.ts b/src/worker/vm/transforms/loop-timeout.ts similarity index 100% rename from src/vm/transforms/loop-timeout.ts rename to src/worker/vm/transforms/loop-timeout.ts diff --git a/src/vm/transforms/promise-timeout.ts b/src/worker/vm/transforms/promise-timeout.ts similarity index 100% rename from src/vm/transforms/promise-timeout.ts rename to src/worker/vm/transforms/promise-timeout.ts diff --git a/src/vm/vm.ts b/src/worker/vm/vm.ts similarity index 99% rename from src/vm/vm.ts rename to src/worker/vm/vm.ts index 6ff353c3..25100f0e 100644 --- a/src/vm/vm.ts +++ b/src/worker/vm/vm.ts @@ -2,7 +2,7 @@ import { randomBytes } from 'crypto' import fetch from 'node-fetch' import { VM } from 'vm2' -import { PluginConfig, PluginConfigVMReponse, PluginsServer } from '../types' +import { PluginConfig, PluginConfigVMReponse, PluginsServer } from '../../types' import { createCache } from './extensions/cache' import { createConsole } from './extensions/console' import { createGoogle } from './extensions/google' diff --git a/src/worker/worker.ts b/src/worker/worker.ts index ea4c3071..38944d7b 100644 --- a/src/worker/worker.ts +++ b/src/worker/worker.ts @@ -1,11 +1,11 @@ -import { ingestEvent } from '../ingestion/ingest-event' import { initApp } from '../init' -import { runPlugins, runPluginsOnBatch, runPluginTask } from '../plugins/run' -import { loadSchedule, setupPlugins } from '../plugins/setup' -import { createServer } from '../server' -import { status } from '../status' +import { createServer } from '../shared/server' +import { status } from '../shared/status' +import { cloneObject } from '../shared/utils' import { PluginsServer, PluginsServerConfig } from '../types' -import { cloneObject } from '../utils' +import { ingestEvent } from './ingestion/ingest-event' +import { runPlugins, runPluginsOnBatch, runPluginTask } from './plugins/run' +import { loadSchedule, setupPlugins } from './plugins/setup' type TaskWorker = ({ task, args }: { task: string; args: any }) => Promise diff --git a/tests/clickhouse/e2e.test.ts b/tests/clickhouse/e2e.test.ts index 737c8b77..318718e7 100644 --- a/tests/clickhouse/e2e.test.ts +++ b/tests/clickhouse/e2e.test.ts @@ -1,10 +1,10 @@ -import { KAFKA_EVENTS_PLUGIN_INGESTION } from '../../src/ingestion/topics' -import { startPluginsServer } from '../../src/server' +import { startPluginsServer } from '../../src/main/pluginsServer' +import { KAFKA_EVENTS_PLUGIN_INGESTION } from '../../src/shared/ingestion/topics' +import { delay, UUIDT } from '../../src/shared/utils' import { LogLevel, PluginsServerConfig } from '../../src/types' import { PluginsServer } from '../../src/types' -import { delay, UUIDT } from '../../src/utils' -import { createPosthog, DummyPostHog } from '../../src/vm/extensions/posthog' import { makePiscina } from '../../src/worker/piscina' +import { createPosthog, DummyPostHog } from '../../src/worker/vm/extensions/posthog' import { resetTestDatabaseClickhouse } from '../helpers/clickhouse' import { resetKafka } from '../helpers/kafka' import { pluginConfig39 } from '../helpers/plugins' diff --git a/tests/clickhouse/ingestion-utils.test.ts b/tests/clickhouse/ingestion-utils.test.ts index 4227ab4a..66816d9a 100644 --- a/tests/clickhouse/ingestion-utils.test.ts +++ b/tests/clickhouse/ingestion-utils.test.ts @@ -1,4 +1,4 @@ -import { chainToElements, elementsToString } from '../../src/ingestion/utils' +import { chainToElements, elementsToString } from '../../src/shared/ingestion/utils' test('elementsToString and chainToElements', () => { const elementsString = elementsToString([ diff --git a/tests/clickhouse/postgres-parity.test.ts b/tests/clickhouse/postgres-parity.test.ts index 9844ab37..3e8ee533 100644 --- a/tests/clickhouse/postgres-parity.test.ts +++ b/tests/clickhouse/postgres-parity.test.ts @@ -1,10 +1,10 @@ import { DateTime } from 'luxon' -import { startPluginsServer } from '../../src/server' +import { startPluginsServer } from '../../src/main/pluginsServer' +import { castTimestampOrNow, UUIDT } from '../../src/shared/utils' import { Database, LogLevel, PluginsServer, PluginsServerConfig, Team, TimestampFormat } from '../../src/types' -import { castTimestampOrNow, UUIDT } from '../../src/utils' -import { createPosthog, DummyPostHog } from '../../src/vm/extensions/posthog' import { makePiscina } from '../../src/worker/piscina' +import { createPosthog, DummyPostHog } from '../../src/worker/vm/extensions/posthog' import { resetTestDatabaseClickhouse } from '../helpers/clickhouse' import { resetKafka } from '../helpers/kafka' import { pluginConfig39 } from '../helpers/plugins' diff --git a/tests/clickhouse/process-event.test.ts b/tests/clickhouse/process-event.test.ts index 8d3e5068..e7175ec2 100644 --- a/tests/clickhouse/process-event.test.ts +++ b/tests/clickhouse/process-event.test.ts @@ -1,4 +1,4 @@ -import { KAFKA_EVENTS_PLUGIN_INGESTION } from '../../src/ingestion/topics' +import { KAFKA_EVENTS_PLUGIN_INGESTION } from '../../src/shared/ingestion/topics' import { Event, PluginsServerConfig } from '../../src/types' import { resetTestDatabaseClickhouse } from '../helpers/clickhouse' import { resetKafka } from '../helpers/kafka' diff --git a/tests/config.test.ts b/tests/config.test.ts index 19ab55e5..63749104 100644 --- a/tests/config.test.ts +++ b/tests/config.test.ts @@ -1,4 +1,4 @@ -import { getDefaultConfig, overrideWithEnv } from '../src/config' +import { getDefaultConfig, overrideWithEnv } from '../src/shared/config' test('overrideWithEnv 1', () => { const defaultConfig = getDefaultConfig() diff --git a/tests/helpers/clickhouse.ts b/tests/helpers/clickhouse.ts index 33af10a8..35c4a5bc 100644 --- a/tests/helpers/clickhouse.ts +++ b/tests/helpers/clickhouse.ts @@ -1,6 +1,6 @@ import ClickHouse from '@posthog/clickhouse' -import { defaultConfig } from '../../src/config' +import { defaultConfig } from '../../src/shared/config' import { PluginsServerConfig } from '../../src/types' export async function resetTestDatabaseClickhouse(extraServerConfig: Partial): Promise { diff --git a/tests/helpers/kafka.ts b/tests/helpers/kafka.ts index 5d110e75..20b60936 100644 --- a/tests/helpers/kafka.ts +++ b/tests/helpers/kafka.ts @@ -1,6 +1,6 @@ import { Kafka, logLevel } from 'kafkajs' -import { defaultConfig, overrideWithEnv } from '../../src/config' +import { defaultConfig, overrideWithEnv } from '../../src/shared/config' import { KAFKA_EVENTS, KAFKA_EVENTS_PLUGIN_INGESTION, @@ -8,9 +8,9 @@ import { KAFKA_PERSON, KAFKA_PERSON_UNIQUE_ID, KAFKA_SESSION_RECORDING_EVENTS, -} from '../../src/ingestion/topics' +} from '../../src/shared/ingestion/topics' +import { delay, UUIDT } from '../../src/shared/utils' import { PluginsServerConfig } from '../../src/types' -import { delay, UUIDT } from '../../src/utils' /** Clear the kafka queue */ export async function resetKafka(extraServerConfig: Partial, delayMs = 2000): Promise { diff --git a/tests/helpers/sql.ts b/tests/helpers/sql.ts index 140d7d05..64a330d9 100644 --- a/tests/helpers/sql.ts +++ b/tests/helpers/sql.ts @@ -1,8 +1,8 @@ import { Pool, PoolClient } from 'pg' -import { defaultConfig } from '../../src/config' +import { defaultConfig } from '../../src/shared/config' +import { delay, UUIDT } from '../../src/shared/utils' import { PluginsServer, PluginsServerConfig, Team } from '../../src/types' -import { delay, UUIDT } from '../../src/utils' import { commonOrganizationId, commonOrganizationMembershipId, commonUserId, makePluginObjects } from './plugins' export async function resetTestDatabase( diff --git a/tests/helpers/sqlMock.ts b/tests/helpers/sqlMock.ts index 42140942..35392325 100644 --- a/tests/helpers/sqlMock.ts +++ b/tests/helpers/sqlMock.ts @@ -1,4 +1,4 @@ -import * as s from '../../src/sql' +import * as s from '../../src/shared/sql' // mock functions that get data from postgres and give them the right types type UnPromisify = F extends (...args: infer A) => Promise ? (...args: A) => T : never diff --git a/tests/helpers/worker.ts b/tests/helpers/worker.ts index a12de0c2..b9205df6 100644 --- a/tests/helpers/worker.ts +++ b/tests/helpers/worker.ts @@ -1,6 +1,6 @@ import Piscina from '@posthog/piscina' -import { defaultConfig } from '../../src/config' +import { defaultConfig } from '../../src/shared/config' import { LogLevel } from '../../src/types' import { makePiscina } from '../../src/worker/piscina' diff --git a/tests/plugins.test.ts b/tests/plugins.test.ts index 18f72260..11a4fb17 100644 --- a/tests/plugins.test.ts +++ b/tests/plugins.test.ts @@ -1,12 +1,12 @@ import { PluginEvent } from '@posthog/plugin-scaffold/src/types' import { mocked } from 'ts-jest/utils' -import { clearError, processError } from '../src/error' -import { loadPlugin } from '../src/plugins/loadPlugin' -import { runPlugins } from '../src/plugins/run' -import { loadSchedule, setupPlugins } from '../src/plugins/setup' -import { createServer } from '../src/server' +import { clearError, processError } from '../src/shared/error' +import { createServer } from '../src/shared/server' import { LogLevel, PluginsServer } from '../src/types' +import { loadPlugin } from '../src/worker/plugins/loadPlugin' +import { runPlugins } from '../src/worker/plugins/run' +import { loadSchedule, setupPlugins } from '../src/worker/plugins/setup' import { commonOrganizationId, mockPluginTempFolder, @@ -17,11 +17,11 @@ import { } from './helpers/plugins' import { getPluginAttachmentRows, getPluginConfigRows, getPluginRows, setError } from './helpers/sqlMock' -jest.mock('../src/sql') -jest.mock('../src/status') -jest.mock('../src/error') -jest.mock('../src/plugins/loadPlugin', () => { - const { loadPlugin } = jest.requireActual('../src/plugins/loadPlugin') +jest.mock('../src/shared/sql') +jest.mock('../src/shared/status') +jest.mock('../src/shared/error') +jest.mock('../src/worker/plugins/loadPlugin', () => { + const { loadPlugin } = jest.requireActual('../src/worker/plugins/loadPlugin') return { loadPlugin: jest.fn().mockImplementation(loadPlugin) } }) diff --git a/tests/postgres/e2e.test.ts b/tests/postgres/e2e.test.ts index 0dd83233..f7c989c8 100644 --- a/tests/postgres/e2e.test.ts +++ b/tests/postgres/e2e.test.ts @@ -1,16 +1,16 @@ import * as IORedis from 'ioredis' -import { startPluginsServer } from '../../src/server' +import { startPluginsServer } from '../../src/main/pluginsServer' +import { UUIDT } from '../../src/shared/utils' import { LogLevel } from '../../src/types' import { PluginsServer } from '../../src/types' -import { UUIDT } from '../../src/utils' -import { createPosthog, DummyPostHog } from '../../src/vm/extensions/posthog' import { makePiscina } from '../../src/worker/piscina' +import { createPosthog, DummyPostHog } from '../../src/worker/vm/extensions/posthog' import { pluginConfig39 } from '../helpers/plugins' import { resetTestDatabase } from '../helpers/sql' import { delayUntilEventIngested } from '../shared/process-event' -jest.mock('../../src/status') +jest.mock('../../src/shared/status') jest.setTimeout(60000) // 60 sec timeout describe('e2e postgres ingestion', () => { diff --git a/tests/postgres/e2e.timeout.test.ts b/tests/postgres/e2e.timeout.test.ts index 23152002..30923a0f 100644 --- a/tests/postgres/e2e.timeout.test.ts +++ b/tests/postgres/e2e.timeout.test.ts @@ -1,8 +1,8 @@ -import { startPluginsServer } from '../../src/server' +import { startPluginsServer } from '../../src/main/pluginsServer' +import { UUIDT } from '../../src/shared/utils' import { LogLevel, PluginsServer } from '../../src/types' -import { UUIDT } from '../../src/utils' -import { createPosthog, DummyPostHog } from '../../src/vm/extensions/posthog' import { makePiscina } from '../../src/worker/piscina' +import { createPosthog, DummyPostHog } from '../../src/worker/vm/extensions/posthog' import { pluginConfig39 } from '../helpers/plugins' import { resetTestDatabase } from '../helpers/sql' import { delayUntilEventIngested } from '../shared/process-event' diff --git a/tests/postgres/queue.test.ts b/tests/postgres/queue.test.ts index a81e2a43..a28b43c3 100644 --- a/tests/postgres/queue.test.ts +++ b/tests/postgres/queue.test.ts @@ -1,9 +1,9 @@ -import Client from '../../src/celery/client' -import { runPlugins } from '../../src/plugins/run' -import { createServer } from '../../src/server' +import { startQueue } from '../../src/main/queue' +import Client from '../../src/shared/celery/client' +import { createServer } from '../../src/shared/server' +import { delay } from '../../src/shared/utils' import { LogLevel, PluginsServer } from '../../src/types' -import { delay } from '../../src/utils' -import { startQueue } from '../../src/worker/queue' +import { runPlugins } from '../../src/worker/plugins/run' jest.setTimeout(60000) // 60 sec timeout diff --git a/tests/postgres/vm.lazy.test.ts b/tests/postgres/vm.lazy.test.ts index 00ccfea5..5435652a 100644 --- a/tests/postgres/vm.lazy.test.ts +++ b/tests/postgres/vm.lazy.test.ts @@ -1,13 +1,13 @@ import { mocked } from 'ts-jest/utils' -import { clearError, processError } from '../../src/error' -import { status } from '../../src/status' -import { LazyPluginVM } from '../../src/vm/lazy' -import { createPluginConfigVM } from '../../src/vm/vm' +import { clearError, processError } from '../../src/shared/error' +import { status } from '../../src/shared/status' +import { LazyPluginVM } from '../../src/worker/vm/lazy' +import { createPluginConfigVM } from '../../src/worker/vm/vm' -jest.mock('../../src/vm/vm') -jest.mock('../../src/error') -jest.mock('../../src/status') +jest.mock('../../src/worker/vm/vm') +jest.mock('../../src/shared/error') +jest.mock('../../src/shared/status') describe('LazyPluginVM', () => { const createVM = () => new LazyPluginVM() diff --git a/tests/postgres/vm.test.ts b/tests/postgres/vm.test.ts index 75d78ebb..d5ca1e70 100644 --- a/tests/postgres/vm.test.ts +++ b/tests/postgres/vm.test.ts @@ -1,15 +1,15 @@ import { PluginEvent } from '@posthog/plugin-scaffold' import * as fetch from 'node-fetch' -import Client from '../../src/celery/client' -import { createServer } from '../../src/server' +import Client from '../../src/shared/celery/client' +import { createServer } from '../../src/shared/server' +import { delay } from '../../src/shared/utils' import { PluginsServer } from '../../src/types' -import { delay } from '../../src/utils' -import { createPluginConfigVM } from '../../src/vm/vm' +import { createPluginConfigVM } from '../../src/worker/vm/vm' import { pluginConfig39 } from '../helpers/plugins' import { resetTestDatabase } from '../helpers/sql' -jest.mock('../../src/celery/client') +jest.mock('../../src/shared/celery/client') const defaultEvent = { distinct_id: 'my_id', diff --git a/tests/postgres/vm.timeout.test.ts b/tests/postgres/vm.timeout.test.ts index 996e250d..1530333c 100644 --- a/tests/postgres/vm.timeout.test.ts +++ b/tests/postgres/vm.timeout.test.ts @@ -1,6 +1,6 @@ -import { createServer } from '../../src/server' +import { createServer } from '../../src/shared/server' import { PluginsServer } from '../../src/types' -import { createPluginConfigVM } from '../../src/vm/vm' +import { createPluginConfigVM } from '../../src/worker/vm/vm' import { pluginConfig39 } from '../helpers/plugins' import { resetTestDatabase } from '../helpers/sql' diff --git a/tests/postgres/worker.test.ts b/tests/postgres/worker.test.ts index 6bc3cbcf..3cce72f1 100644 --- a/tests/postgres/worker.test.ts +++ b/tests/postgres/worker.test.ts @@ -2,24 +2,24 @@ import { PluginEvent } from '@posthog/plugin-scaffold/src/types' import IORedis from 'ioredis' import { mocked } from 'ts-jest/utils' -import Client from '../../src/celery/client' -import { ingestEvent } from '../../src/ingestion/ingest-event' -import { runPlugins, runPluginsOnBatch, runPluginTask } from '../../src/plugins/run' -import { loadSchedule, setupPlugins } from '../../src/plugins/setup' -import { ServerInstance, startPluginsServer } from '../../src/server' -import { loadPluginSchedule } from '../../src/services/schedule' +import { ServerInstance, startPluginsServer } from '../../src/main/pluginsServer' +import { loadPluginSchedule } from '../../src/main/services/schedule' +import Client from '../../src/shared/celery/client' +import { delay, UUIDT } from '../../src/shared/utils' import { LogLevel } from '../../src/types' -import { delay, UUIDT } from '../../src/utils' +import { ingestEvent } from '../../src/worker/ingestion/ingest-event' import { makePiscina } from '../../src/worker/piscina' +import { runPlugins, runPluginsOnBatch, runPluginTask } from '../../src/worker/plugins/run' +import { loadSchedule, setupPlugins } from '../../src/worker/plugins/setup' import { createTaskRunner } from '../../src/worker/worker' import { resetTestDatabase } from '../helpers/sql' import { setupPiscina } from '../helpers/worker' -jest.mock('../../src/sql') -jest.mock('../../src/status') -jest.mock('../../src/ingestion/ingest-event') -jest.mock('../../src/plugins/run') -jest.mock('../../src/plugins/setup') +jest.mock('../../src/shared/sql') +jest.mock('../../src/shared/status') +jest.mock('../../src/worker/ingestion/ingest-event') +jest.mock('../../src/worker/plugins/run') +jest.mock('../../src/worker/plugins/setup') jest.setTimeout(600000) // 600 sec timeout function createEvent(index = 0): PluginEvent { @@ -92,7 +92,7 @@ test('assume that the workerThreads and tasksPerWorker values behave as expected await resetTestDatabase(testCode) const piscina = setupPiscina(workerThreads, tasksPerWorker) const processEvent = (event: PluginEvent) => piscina.runTask({ task: 'processEvent', args: { event } }) - const promises = [] + const promises: Array> = [] // warmup 2x await Promise.all([processEvent(createEvent()), processEvent(createEvent())]) diff --git a/tests/schedule.test.ts b/tests/schedule.test.ts index a4116e97..e6fabb72 100644 --- a/tests/schedule.test.ts +++ b/tests/schedule.test.ts @@ -1,21 +1,21 @@ import { PluginEvent } from '@posthog/plugin-scaffold/src/types' -import { createServer } from '../src/server' import { loadPluginSchedule, LOCKED_RESOURCE, runTasksDebounced, startSchedule, waitForTasksToFinish, -} from '../src/services/schedule' +} from '../src/main/services/schedule' +import { createServer } from '../src/shared/server' +import { delay } from '../src/shared/utils' import { LogLevel, ScheduleControl } from '../src/types' -import { delay } from '../src/utils' import { createPromise } from './helpers/promises' import { resetTestDatabase } from './helpers/sql' import { setupPiscina } from './helpers/worker' -jest.mock('../src/sql') -jest.mock('../src/status') +jest.mock('../src/shared/sql') +jest.mock('../src/shared/status') jest.setTimeout(60000) // 60 sec timeout function createEvent(index = 0): PluginEvent { diff --git a/tests/server.test.ts b/tests/server.test.ts index b29b1b5e..f6616beb 100644 --- a/tests/server.test.ts +++ b/tests/server.test.ts @@ -1,11 +1,9 @@ -import { PluginEvent } from '@posthog/plugin-scaffold/src/types' - -import { startPluginsServer } from '../src/server' +import { startPluginsServer } from '../src/main/pluginsServer' import { LogLevel } from '../src/types' import { makePiscina } from '../src/worker/piscina' import { resetTestDatabase } from './helpers/sql' -jest.mock('../src/sql') +jest.mock('../src/shared/sql') jest.setTimeout(60000) // 60 sec timeout test('startPluginsServer', async () => { diff --git a/tests/shared/process-event.ts b/tests/shared/process-event.ts index 02a74f58..0dcd12d1 100644 --- a/tests/shared/process-event.ts +++ b/tests/shared/process-event.ts @@ -4,9 +4,9 @@ import { DateTime } from 'luxon' import { performance } from 'perf_hooks' import { IEvent } from '../../src/idl/protos' -import { EventsProcessor } from '../../src/ingestion/process-event' -import { hashElements } from '../../src/ingestion/utils' -import { createServer } from '../../src/server' +import { hashElements } from '../../src/shared/ingestion/utils' +import { createServer } from '../../src/shared/server' +import { delay, UUIDT } from '../../src/shared/utils' import { Database, Event, @@ -17,7 +17,7 @@ import { SessionRecordingEvent, Team, } from '../../src/types' -import { delay, UUIDT } from '../../src/utils' +import { EventsProcessor } from '../../src/worker/ingestion/process-event' import { createUserTeamAndOrganization, getFirstTeam, getTeams, onQuery, resetTestDatabase } from '../helpers/sql' jest.setTimeout(600000) // 600 sec timeout. diff --git a/tests/sql.test.ts b/tests/sql.test.ts index 171da780..571fbc2f 100644 --- a/tests/sql.test.ts +++ b/tests/sql.test.ts @@ -1,5 +1,5 @@ -import { createServer } from '../src/server' -import { getPluginAttachmentRows, getPluginConfigRows, getPluginRows, setError } from '../src/sql' +import { createServer } from '../src/shared/server' +import { getPluginAttachmentRows, getPluginConfigRows, getPluginRows, setError } from '../src/shared/sql' import { PluginConfig, PluginError, PluginsServer } from '../src/types' import { commonOrganizationId } from './helpers/plugins' import { resetTestDatabase } from './helpers/sql' diff --git a/tests/transforms.test.ts b/tests/transforms.test.ts index ed4aa3a1..726fe7fb 100644 --- a/tests/transforms.test.ts +++ b/tests/transforms.test.ts @@ -1,7 +1,7 @@ -import { createServer } from '../src/server' +import { createServer } from '../src/shared/server' +import { code } from '../src/shared/utils' import { PluginsServer } from '../src/types' -import { code } from '../src/utils' -import { transformCode } from '../src/vm/transforms' +import { transformCode } from '../src/worker/vm/transforms' import { resetTestDatabase } from './helpers/sql' let server: PluginsServer diff --git a/tests/utils.test.ts b/tests/utils.test.ts index ac397f48..d9ba4e2e 100644 --- a/tests/utils.test.ts +++ b/tests/utils.test.ts @@ -1,6 +1,5 @@ import { randomBytes } from 'crypto' -import { LogLevel } from '../src/types' import { bufferToStream, cloneObject, @@ -12,7 +11,8 @@ import { setLogLevel, UUID, UUIDT, -} from '../src/utils' +} from '../src/shared/utils' +import { LogLevel } from '../src/types' // .zip in Base64: github repo posthog/helloworldplugin const zip =