From 60859c04da8f74b08578eb96f4ce1e04777f62c6 Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Thu, 4 Feb 2021 13:44:01 +0100 Subject: [PATCH 1/2] merge test and postgres query simplification --- src/ingestion/process-event.ts | 15 +++------ tests/shared/process-event.ts | 57 ++++++++++++++++++++++++++++++---- 2 files changed, 56 insertions(+), 16 deletions(-) diff --git a/src/ingestion/process-event.ts b/src/ingestion/process-event.ts index acc584e6..06657bdb 100644 --- a/src/ingestion/process-event.ts +++ b/src/ingestion/process-event.ts @@ -265,7 +265,7 @@ export class EventsProcessor { } } - private async mergePeople(mergeInto: Person, peopleToMerge: Person[]): Promise { + public async mergePeople(mergeInto: Person, peopleToMerge: Person[]): Promise { let firstSeen = mergeInto.created_at // merge the properties @@ -291,15 +291,10 @@ export class EventsProcessor { await this.db.moveDistinctId(otherPerson, personDistinctId, mergeInto) } - const otherCohortPeople: CohortPeople[] = ( - await this.db.postgresQuery('SELECT * FROM posthog_cohortpeople WHERE person_id = $1', [otherPerson.id]) - ).rows - for (const cohortPeople of otherCohortPeople) { - await this.db.postgresQuery('UPDATE posthog_cohortpeople SET person_id = $1 WHERE id = $2', [ - mergeInto.id, - cohortPeople.id, - ]) - } + await this.db.postgresQuery('UPDATE posthog_cohortpeople SET person_id = $1 WHERE person_id = $2', [ + mergeInto.id, + otherPerson.id, + ]) await this.db.deletePerson(otherPerson) } diff --git a/tests/shared/process-event.ts b/tests/shared/process-event.ts index 249b67ab..c38eb59e 100644 --- a/tests/shared/process-event.ts +++ b/tests/shared/process-event.ts @@ -1,16 +1,14 @@ import { PluginEvent } from '@posthog/plugin-scaffold/src/types' import { createServer } from '../../src/server' import { - LogLevel, - PluginsServer, - Team, + Database, Event, + LogLevel, Person, - Element, - PostgresSessionRecordingEvent, + PluginsServer, PluginsServerConfig, - ClickHouseEvent, SessionRecordingEvent, + Team, } from '../../src/types' import { createUserTeamAndOrganization, getFirstTeam, getTeams, resetTestDatabase } from '../helpers/sql' import { EventsProcessor } from '../../src/ingestion/process-event' @@ -125,6 +123,53 @@ export const createProcessEventTests = ( createTests?.(returned) + test('merge people', async () => { + const p0 = await createPerson(server, team, ['person_0'], { $os: 'Microsoft' }) + await delayUntilEventIngested(() => server.db.fetchPersons(Database.ClickHouse), 1) + await server.db.updatePerson(p0, { created_at: DateTime.fromISO('2020-01-01T00:00:00Z') }) + + const p1 = await createPerson(server, team, ['person_1'], { $os: 'Chrome' }) + await delayUntilEventIngested(() => server.db.fetchPersons(Database.ClickHouse), 2) + await server.db.updatePerson(p1, { created_at: DateTime.fromISO('2019-07-01T00:00:00Z') }) + + await processEvent( + 'person_1', + '', + '', + ({ + event: 'user signed up', + properties: {}, + } as any) as PluginEvent, + team.id, + now, + now, + new UUIDT().toString() + ) + + await createPerson(server, team, ['person_2'], { $os: 'Apple', $browser: 'MS Edge' }) + await createPerson(server, team, ['person_3'], { $os: 'PlayStation' }) + + await delayUntilEventIngested(() => server.db.fetchPersons(Database.ClickHouse), 4) + + const [person0, person1, person2, person3] = await server.db.fetchPersons() + + expect((await server.db.fetchPersons(Database.ClickHouse)).length).toEqual(4) + + await eventsProcessor.mergePeople(person0, [person1, person2, person3]) + + await delayUntilEventIngested(async () => + (await server.db.fetchPersons(Database.ClickHouse)).length === 1 ? [1] : [] + ) + + expect((await server.db.fetchPersons(Database.ClickHouse)).length).toEqual(1) + + const [person] = await server.db.fetchPersons() + + expect(person.properties).toEqual({ $os: 'Microsoft', $browser: 'MS Edge' }) + expect(await server.db.fetchDistinctIdValues(person)).toEqual(['person_0', 'person_1', 'person_2', 'person_3']) + expect(person.created_at.toISO()).toEqual(DateTime.fromISO('2019-07-01T00:00:00Z').setZone('UTC').toISO()) + }) + test('capture new person', async () => { await server.db.postgresQuery(`UPDATE posthog_team SET ingested_event = $1 WHERE id = $2`, [true, team.id]) team = await getFirstTeam(server) From d73f4f4c6c010cd158f921bd7cd57c6d56df4054 Mon Sep 17 00:00:00 2001 From: Marius Andra Date: Thu, 4 Feb 2021 13:47:20 +0100 Subject: [PATCH 2/2] postgres fix --- tests/shared/process-event.ts | 28 +++++++++++++++++++--------- 1 file changed, 19 insertions(+), 9 deletions(-) diff --git a/tests/shared/process-event.ts b/tests/shared/process-event.ts index c38eb59e..b5e8903e 100644 --- a/tests/shared/process-event.ts +++ b/tests/shared/process-event.ts @@ -125,11 +125,16 @@ export const createProcessEventTests = ( test('merge people', async () => { const p0 = await createPerson(server, team, ['person_0'], { $os: 'Microsoft' }) - await delayUntilEventIngested(() => server.db.fetchPersons(Database.ClickHouse), 1) + if (database === 'clickhouse') { + await delayUntilEventIngested(() => server.db.fetchPersons(Database.ClickHouse), 1) + } + await server.db.updatePerson(p0, { created_at: DateTime.fromISO('2020-01-01T00:00:00Z') }) const p1 = await createPerson(server, team, ['person_1'], { $os: 'Chrome' }) - await delayUntilEventIngested(() => server.db.fetchPersons(Database.ClickHouse), 2) + if (database === 'clickhouse') { + await delayUntilEventIngested(() => server.db.fetchPersons(Database.ClickHouse), 2) + } await server.db.updatePerson(p1, { created_at: DateTime.fromISO('2019-07-01T00:00:00Z') }) await processEvent( @@ -149,19 +154,24 @@ export const createProcessEventTests = ( await createPerson(server, team, ['person_2'], { $os: 'Apple', $browser: 'MS Edge' }) await createPerson(server, team, ['person_3'], { $os: 'PlayStation' }) - await delayUntilEventIngested(() => server.db.fetchPersons(Database.ClickHouse), 4) + if (database === 'clickhouse') { + await delayUntilEventIngested(() => server.db.fetchPersons(Database.ClickHouse), 4) + expect((await server.db.fetchPersons(Database.ClickHouse)).length).toEqual(4) + } + expect((await server.db.fetchPersons()).length).toEqual(4) const [person0, person1, person2, person3] = await server.db.fetchPersons() - expect((await server.db.fetchPersons(Database.ClickHouse)).length).toEqual(4) - await eventsProcessor.mergePeople(person0, [person1, person2, person3]) - await delayUntilEventIngested(async () => - (await server.db.fetchPersons(Database.ClickHouse)).length === 1 ? [1] : [] - ) + if (database === 'clickhouse') { + await delayUntilEventIngested(async () => + (await server.db.fetchPersons(Database.ClickHouse)).length === 1 ? [1] : [] + ) + expect((await server.db.fetchPersons(Database.ClickHouse)).length).toEqual(1) + } - expect((await server.db.fetchPersons(Database.ClickHouse)).length).toEqual(1) + expect((await server.db.fetchPersons()).length).toEqual(1) const [person] = await server.db.fetchPersons()