Skip to content
Open
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
2 changes: 2 additions & 0 deletions services/apps/script_executor_worker/src/activities.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ import {
import {
getWorkflowsCount,
mergeMembers,
mergeMembersIfAllowed,
mergeOrganizations,
triggerMemberAffiliationsRefresh,
unmergeMembers,
Expand Down Expand Up @@ -66,6 +67,7 @@ export {
findMembersWithSameVerifiedEmailsInDifferentPlatforms,
findMembersWithSamePlatformIdentitiesDifferentCapitalization,
mergeMembers,
mergeMembersIfAllowed,
findMemberMergeActions,
findMergeActionUnmergeBackup,
unmergeMembers,
Expand Down
54 changes: 53 additions & 1 deletion services/apps/script_executor_worker/src/activities/common.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,8 @@
import axios from 'axios'

import { CommonMemberService, signalMemberUpdate } from '@crowd/common_services'
import { pgpQx } from '@crowd/data-access-layer'
import { findExistingMemberIds, pgpQx } from '@crowd/data-access-layer'
import { getMemberNoMerge } from '@crowd/data-access-layer/src/member_merge'
import {
IMemberIdentity,
IMemberUnmergeBackup,
Expand All @@ -27,6 +28,57 @@ export async function mergeMembers(
}
}

export async function mergeMembersIfAllowed(
primaryMemberId: string,
secondaryMemberId: string,
): Promise<boolean> {
const qx = pgpQx(svc.postgres.writer.connection())

const existingMemberIds = await findExistingMemberIds(qx, [primaryMemberId, secondaryMemberId])

if (existingMemberIds.length < 2) {
svc.log.info(
{ primaryMemberId, secondaryMemberId },
'One of the members no longer exists - skipping merge!',
)
return false
}

const noMergeIds = await getMemberNoMerge(qx, [primaryMemberId, secondaryMemberId])
const blockedByNoMerge = noMergeIds.some(
(m) =>
(m.memberId === primaryMemberId && m.noMergeId === secondaryMemberId) ||
(m.memberId === secondaryMemberId && m.noMergeId === primaryMemberId),
)
Comment thread
ramanathan1504 marked this conversation as resolved.
Comment thread
ramanathan1504 marked this conversation as resolved.

if (blockedByNoMerge) {
svc.log.warn(
{ primaryMemberId, secondaryMemberId },
'Members are marked as no-merge - skipping merge!',
)
return false
}

const memberService = new CommonMemberService(qx, svc.temporal, svc.log)

try {
await memberService.merge(primaryMemberId, secondaryMemberId)
} catch (error) {
if (error?.code === 409) {
svc.log.info(
{ primaryMemberId, secondaryMemberId },
'Another merge is already in progress - skipping merge!',
)
return false
}

svc.log.error({ err: error, primaryMemberId, secondaryMemberId }, 'Failed to merge members')
throw error
}

return true
}

export async function unmergeMembers(
primaryMemberId: string,
backup: IUnmergeBackup<IMemberUnmergeBackup> | IUnmergePreviewResult<IMemberUnmergePreviewResult>,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,13 +5,18 @@ import { svc } from '../../main'

export async function findMembersWithSameVerifiedEmailsInDifferentPlatforms(
limit: number,
afterHash?: number,
afterHighMemberId?: string,
afterLowMemberId?: string,
): Promise<ISimilarMember[]> {
let rows: ISimilarMember[] = []

try {
const memberRepo = new MemberRepository(svc.postgres.reader.connection(), svc.log)
rows = await memberRepo.findMembersWithSameVerifiedEmailsInDifferentPlatforms(limit, afterHash)
rows = await memberRepo.findMembersWithSameVerifiedEmailsInDifferentPlatforms(
limit,
afterHighMemberId,
afterLowMemberId,
)
} catch (err) {
throw new Error(err)
}
Expand Down
2 changes: 2 additions & 0 deletions services/apps/script_executor_worker/src/main.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import {
scheduleOrganizationSegmentAggCleanup,
scheduleOrganizationsCleanup,
} from './schedules/scheduleCleanup'
import { scheduleMergeMembersWithSameVerifiedEmails } from './schedules/scheduleMemberDeduplication'

const config: Config = {
envvars: [
Expand Down Expand Up @@ -47,6 +48,7 @@ setImmediate(async () => {
await scheduleOrganizationsCleanup()
await scheduleMemberSegmentsAggCleanup()
await scheduleOrganizationSegmentAggCleanup()
await scheduleMergeMembersWithSameVerifiedEmails()

await svc.start()
})
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
import { ScheduleAlreadyRunning, ScheduleOverlapPolicy } from '@temporalio/client'

import { svc } from '../main'
import { findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatforms } from '../workflows/findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatforms'

export const scheduleMergeMembersWithSameVerifiedEmails = async () => {
try {
await svc.temporal.schedule.create({
scheduleId: 'mergeMembersWithSameVerifiedEmails',
spec: {
cronExpressions: ['0 3 * * 0'],
},
policies: {
overlap: ScheduleOverlapPolicy.SKIP,
catchupWindow: '1 minute',
},
action: {
type: 'startWorkflow',
workflowType: findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatforms,
Comment thread
ramanathan1504 marked this conversation as resolved.
taskQueue: 'script-executor',
retry: {
initialInterval: '15 seconds',
backoffCoefficient: 2,
maximumAttempts: 3,
},
args: [{}],
},
})
svc.log.info('Schedule for merging members with same verified emails created successfully!')
} catch (err) {
if (err instanceof ScheduleAlreadyRunning) {
svc.log.info('Schedule mergeMembersWithSameVerifiedEmails already registered in Temporal.')
svc.log.info('Configuration may have changed since. Please make sure they are in sync.')
} else {
svc.log.error({ err }, 'Error creating schedule for member email deduplication')
throw new Error(err)
}
}
}
4 changes: 3 additions & 1 deletion services/apps/script_executor_worker/src/types.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
import { IActivityRelationDuplicateGroup } from '@crowd/data-access-layer'

export interface IFindAndMergeMembersWithSameVerifiedEmailsInDifferentPlatformsArgs {
afterHash?: number
afterHighMemberId?: string
afterLowMemberId?: string
dryRun?: boolean
}

export interface IFindAndMergeMembersWithSameIdentitiesDifferentCapitalizationInPlatformArgs {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,8 @@ export async function findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatfo
const mergeableMemberCouples =
await activity.findMembersWithSameVerifiedEmailsInDifferentPlatforms(
PROCESS_MEMBERS_PER_RUN,
args.afterHash || undefined,
args.afterHighMemberId || undefined,
args.afterLowMemberId || undefined,
)

if (mergeableMemberCouples.length === 0) {
Expand All @@ -31,13 +32,27 @@ export async function findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatfo
}

for (const couple of mergeableMemberCouples) {
console.log(
`Merging ${couple.secondaryMemberId} [${couple.secondaryMemberIdentityValue}] into ${couple.primaryMemberId} [${couple.primaryMemberIdentityValue}]! `,
)
await common.mergeMembers(couple.primaryMemberId, couple.secondaryMemberId)
const coupleDescription = `${couple.secondaryMemberId} [${couple.secondaryMemberIdentityValue}] into ${couple.primaryMemberId} [${couple.primaryMemberIdentityValue}]`

if (args.dryRun) {
console.log(`[dry run] Would merge ${coupleDescription}!`)
continue
}

console.log(`Merging ${coupleDescription}!`)

await common.mergeMembersIfAllowed(couple.primaryMemberId, couple.secondaryMemberId)
}

const lastCouple = mergeableMemberCouples[mergeableMemberCouples.length - 1]
const [afterLowMemberId, afterHighMemberId] = [
lastCouple.primaryMemberId,
lastCouple.secondaryMemberId,
].sort()

await continueAsNew<typeof findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatforms>({
afterHash: mergeableMemberCouples[mergeableMemberCouples.length - 1]?.hash,
afterHighMemberId,
afterLowMemberId,
dryRun: args.dryRun,
})
}
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@ import {
preferCompanyOverUniversityWhenOverlapping,
updateMember,
} from '@crowd/data-access-layer'
import { removeMemberToMerge } from '@crowd/data-access-layer/src/member_merge'
import { moveMemberNoMerge, removeMemberToMerge } from '@crowd/data-access-layer/src/member_merge'
import {
deleteMemberSegmentAffiliations,
findMemberAffiliations,
Expand Down Expand Up @@ -421,6 +421,8 @@ export class CommonMemberService extends LoggerBase {
identitiesToUpdate,
)

await moveMemberNoMerge(txQx, toMergeId, originalId)

// Update member segment affiliations and organization affiliation overrides
await moveAffiliationsBetweenMembers(txQx, toMergeId, originalId)

Expand Down
27 changes: 27 additions & 0 deletions services/libs/data-access-layer/src/member_merge/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,33 @@ export async function insertMemberNoMerge(
)
}

export async function moveMemberNoMerge(
qx: QueryExecutor,
fromMemberId: string,
toMemberId: string,
): Promise<void> {
await qx.result(
`
with "blockedMembers" as (
select distinct
case when "memberId" = $(fromMemberId) then "noMergeId" else "memberId" end as id
from "memberNoMerge"
where "memberId" = $(fromMemberId) or "noMergeId" = $(fromMemberId)
)
insert into "memberNoMerge" ("memberId", "noMergeId", "createdAt", "updatedAt")
select "memberId", "noMergeId", NOW(), NOW()
from (
select $(toMemberId)::uuid as "memberId", b.id as "noMergeId" from "blockedMembers" b
union
select b.id as "memberId", $(toMemberId)::uuid as "noMergeId" from "blockedMembers" b
) edges
where "memberId" != "noMergeId"
ON CONFLICT ("memberId", "noMergeId") DO NOTHING
`,
{ fromMemberId, toMemberId },
)
}

export async function getMemberNoMerge(
qx: QueryExecutor,
memberIds: string[],
Expand Down
11 changes: 11 additions & 0 deletions services/libs/data-access-layer/src/members/base.ts
Original file line number Diff line number Diff line change
Expand Up @@ -716,6 +716,17 @@ export async function findMemberById<T extends MemberField>(
return queryTableById(qx, 'members', Object.values(MemberField), memberId, fields)
}

export async function findExistingMemberIds(
qx: QueryExecutor,
memberIds: string[],
): Promise<string[]> {
const rows = await qx.select(`select id from members where id in ($(memberIds:csv))`, {
memberIds,
})

return rows.map((row) => row.id)
}

export async function moveAffiliationsBetweenMembers(
qx: QueryExecutor,
fromMemberId: string,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,13 +16,15 @@ class MemberRepository {

async findMembersWithSameVerifiedEmailsInDifferentPlatforms(
limit = 50,
afterHash: number = undefined,
afterHighMemberId: string = undefined,
afterLowMemberId: string = undefined,
): Promise<ISimilarMember[]> {
let rows: ISimilarMember[] = []
try {
const afterHashFilter = afterHash
? ` and Greatest(Hashtext(Concat(a."memberId", b."memberId")), Hashtext(Concat(b."memberId", a."memberId"))) < $(afterHash) `
: ''
const afterPairFilter =
afterHighMemberId && afterLowMemberId
? ` and (Greatest(a."memberId"::text, b."memberId"::text), Least(a."memberId"::text, b."memberId"::text)) < ($(afterHighMemberId), $(afterLowMemberId)) `
: ''
rows = await this.connection.query(
`
select
Expand All @@ -39,13 +41,19 @@ class MemberRepository {
and a.type = 'email'
and a."deletedAt" is null
and b."deletedAt" is null
${afterHashFilter}
group by hash
order by hash desc
${afterPairFilter}
group by
Greatest(a."memberId"::text, b."memberId"::text),
Least(a."memberId"::text, b."memberId"::text),
hash
Comment thread
cursor[bot] marked this conversation as resolved.
order by
Greatest(a."memberId"::text, b."memberId"::text) desc,
Least(a."memberId"::text, b."memberId"::text) desc
limit $(limit);
`,
{
afterHash,
afterHighMemberId,
afterLowMemberId,
limit,
},
)
Expand Down