From cefc84c6b942be2db371d1a201678525f03a5946 Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Mon, 1 Feb 2021 11:29:09 +0100 Subject: [PATCH] postgres e2e ingestion test --- src/config.ts | 2 + src/types.ts | 1 + src/worker/queue.ts | 33 ++++++++++----- tests/clickhouse/process-event.test.ts | 10 +---- tests/postgres/e2e.test.ts | 56 ++++++++++++++++++++++++++ 5 files changed, 84 insertions(+), 18 deletions(-) create mode 100644 tests/postgres/e2e.test.ts diff --git a/src/config.ts b/src/config.ts index 1b6c108f..d3e08dee 100644 --- a/src/config.ts +++ b/src/config.ts @@ -20,6 +20,7 @@ export function getDefaultConfig(): PluginsServerConfig { KAFKA_CLIENT_CERT_B64: null, KAFKA_CLIENT_CERT_KEY_B64: null, KAFKA_TRUSTED_CERT_B64: null, + PLUGIN_SERVER_INGESTION: false, PLUGINS_CELERY_QUEUE: 'posthog-plugins', REDIS_URL: 'redis://127.0.0.1', BASE_DIR: '.', @@ -40,6 +41,7 @@ export function getDefaultConfig(): PluginsServerConfig { export function getConfigHelp(): Record { return { + PLUGIN_SERVER_INGESTION: 'Ingest events via plugin-server', CELERY_DEFAULT_QUEUE: 'Celery outgoing queue', PLUGINS_CELERY_QUEUE: 'Celery incoming queue', DATABASE_URL: 'Postgres database URL', diff --git a/src/types.ts b/src/types.ts index e457f674..5aeb9416 100644 --- a/src/types.ts +++ b/src/types.ts @@ -42,6 +42,7 @@ export interface PluginsServerConfig extends Record { WEB_PORT: number WEB_HOSTNAME: string LOG_LEVEL: LogLevel + PLUGIN_SERVER_INGESTION: boolean SENTRY_DSN: string | null STATSD_HOST: string | null STATSD_PORT: number diff --git a/src/worker/queue.ts b/src/worker/queue.ts index 99ae5ed2..f55f1995 100644 --- a/src/worker/queue.ts +++ b/src/worker/queue.ts @@ -6,6 +6,7 @@ import Client from '../celery/client' import { PluginsServer, Queue } from '../types' import { KafkaQueue } from '../ingestion/kafka-queue' import { status } from '../status' +import { UUIDT } from '../utils' export async function startQueue( server: PluginsServer, @@ -45,15 +46,29 @@ async function startQueueRedis( const processedEvent = await processEvent(event) if (processedEvent) { const { distinct_id, ip, site_url, team_id, now, sent_at, ...data } = processedEvent - client.sendTask('posthog.tasks.process_event.process_event', [], { - distinct_id, - ip, - site_url, - data, - team_id, - now, - sent_at, - }) + + if (server.PLUGIN_SERVER_INGESTION) { + await server.eventsProcessor.processEvent( + distinct_id, + ip, + site_url, + processedEvent, + team_id, + DateTime.fromISO(now), + sent_at ? DateTime.fromISO(sent_at) : null, + new UUIDT().toString() + ) + } else { + client.sendTask('posthog.tasks.process_event.process_event', [], { + distinct_id, + ip, + site_url, + data, + team_id, + now, + sent_at, + }) + } } } catch (e) { Sentry.captureException(e) diff --git a/tests/clickhouse/process-event.test.ts b/tests/clickhouse/process-event.test.ts index b86ed4be..96214920 100644 --- a/tests/clickhouse/process-event.test.ts +++ b/tests/clickhouse/process-event.test.ts @@ -1,12 +1,4 @@ -import { - Element, - Event, - Person, - PersonDistinctId, - PluginsServer, - PluginsServerConfig, - PostgresSessionRecordingEvent, -} from '../../src/types' +import { PluginsServerConfig } from '../../src/types' import { resetTestDatabaseClickhouse } from '../helpers/clickhouse' import { KafkaCollector, KafkaObserver } from '../helpers/kafka' import { UUIDT } from '../../src/utils' diff --git a/tests/postgres/e2e.test.ts b/tests/postgres/e2e.test.ts new file mode 100644 index 00000000..208eef5a --- /dev/null +++ b/tests/postgres/e2e.test.ts @@ -0,0 +1,56 @@ +import { LogLevel } from '../../src/types' +import { resetTestDatabase } from '../helpers/sql' +import { startPluginsServer } from '../../src/server' +import { makePiscina } from '../../src/worker/piscina' +import { PluginsServer } from '../../src/types' +import { createPosthog, DummyPostHog } from '../../src/extensions/posthog' +import { pluginConfig39 } from '../helpers/plugins' +import { delay } from '../../src/utils' + +jest.setTimeout(60000) // 60 sec timeout + +describe('e2e postgres ingestion', () => { + let server: PluginsServer + let stopServer: () => Promise + let posthog: DummyPostHog + + beforeEach(async () => { + await resetTestDatabase(` + async function processEvent (event) { + event.properties.processed = 'hell yes' + return event + } + `) + const startResponse = await startPluginsServer( + { + WORKER_CONCURRENCY: 2, + PLUGINS_CELERY_QUEUE: 'test-plugins-celery-queue', + CELERY_DEFAULT_QUEUE: 'test-celery-default-queue', + PLUGIN_SERVER_INGESTION: true, + LOG_LEVEL: LogLevel.Log, + KAFKA_ENABLED: false, + }, + makePiscina + ) + server = startResponse.server + stopServer = startResponse.stop + + await server.redis.del(server.PLUGINS_CELERY_QUEUE) + await server.redis.del(server.CELERY_DEFAULT_QUEUE) + + posthog = createPosthog(server, pluginConfig39) + }) + + afterEach(async () => { + await stopServer() + }) + + test('event captured, processed, ingested', async () => { + expect((await server.db.fetchEvents()).length).toBe(0) + posthog.capture('custom event', { name: 'haha' }) + await delay(2000) + const events = await server.db.fetchEvents() + expect(events.length).toBe(1) + expect(events[0].properties.processed).toEqual('hell yes') + }) +})