Skip to content
This repository was archived by the owner on Nov 4, 2021. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions benchmarks/clickhouse/e2e.kafka.benchmark.ts
Original file line number Diff line number Diff line change
@@ -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'
Expand Down
8 changes: 4 additions & 4 deletions benchmarks/clickhouse/e2e.timeout.benchmark.ts
Original file line number Diff line number Diff line change
@@ -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'
Expand Down
6 changes: 3 additions & 3 deletions benchmarks/postgres/e2e.celery.benchmark.ts
Original file line number Diff line number Diff line change
@@ -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'
Expand Down
4 changes: 2 additions & 2 deletions benchmarks/postgres/helpers/piscina.ts
Original file line number Diff line number Diff line change
@@ -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 {
Expand Down
8 changes: 4 additions & 4 deletions benchmarks/postgres/ingestion.benchmark.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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', () => {
Expand Down
6 changes: 3 additions & 3 deletions benchmarks/vm/memory.benchmark.ts
Original file line number Diff line number Diff line change
@@ -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 {
Expand Down
4 changes: 2 additions & 2 deletions benchmarks/vm/worker.benchmark.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
4 changes: 2 additions & 2 deletions src/healthcheck.ts
Original file line number Diff line number Diff line change
@@ -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')

Expand Down
4 changes: 2 additions & 2 deletions src/index.ts
Original file line number Diff line number Diff line change
@@ -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'

Expand Down
2 changes: 1 addition & 1 deletion src/init.ts
Original file line number Diff line number Diff line change
@@ -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')
Expand Down
Original file line number Diff line number Diff line change
@@ -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<void>

export class Worker extends Base implements Queue {
export class CeleryQueueWorker extends Base implements Queue {
handlers: Record<string, Handler> = {}
activeTasks: Set<Promise<any>> = new Set()

Expand Down Expand Up @@ -225,4 +225,4 @@ export class Worker extends Base implements Queue {
}
}

export default Worker
export default CeleryQueueWorker
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
149 changes: 149 additions & 0 deletions src/main/pluginsServer.ts
Original file line number Diff line number Diff line change
@@ -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<void>
}

export async function startPluginsServer(
config: Partial<PluginsServerConfig>,
makePiscina: (config: PluginsServerConfig) => Piscina
): Promise<ServerInstance> {
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<void> | undefined
let scheduleControl: ScheduleControl | undefined

let shutdownStatus = 0

async function closeJobs(): Promise<void> {
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<void> {
// 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()
}
15 changes: 7 additions & 8 deletions src/worker/queue.ts → src/main/queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<PluginEvent | null>
Expand Down Expand Up @@ -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(
Expand Down
8 changes: 4 additions & 4 deletions src/services/schedule.ts → src/main/services/schedule.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'

Expand Down
2 changes: 1 addition & 1 deletion src/web/server.ts → src/main/web/server.ts
Original file line number Diff line number Diff line change
@@ -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()
Expand Down
File renamed without changes.
2 changes: 1 addition & 1 deletion src/celery/broker.ts → src/shared/celery/broker.ts
Original file line number Diff line number Diff line change
@@ -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 }
Expand Down
File renamed without changes.
File renamed without changes.
File renamed without changes.
File renamed without changes.
2 changes: 1 addition & 1 deletion src/config.ts → src/shared/config.ts
Original file line number Diff line number Diff line change
@@ -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()
Expand Down
Loading