diff --git a/services/apps/script_executor_worker/src/activities.ts b/services/apps/script_executor_worker/src/activities.ts index 706e5576bb..d89c69530a 100644 --- a/services/apps/script_executor_worker/src/activities.ts +++ b/services/apps/script_executor_worker/src/activities.ts @@ -27,6 +27,7 @@ import { import { getWorkflowsCount, mergeMembers, + mergeMembersIfAllowed, mergeOrganizations, triggerMemberAffiliationsRefresh, unmergeMembers, @@ -66,6 +67,7 @@ export { findMembersWithSameVerifiedEmailsInDifferentPlatforms, findMembersWithSamePlatformIdentitiesDifferentCapitalization, mergeMembers, + mergeMembersIfAllowed, findMemberMergeActions, findMergeActionUnmergeBackup, unmergeMembers, diff --git a/services/apps/script_executor_worker/src/activities/common.ts b/services/apps/script_executor_worker/src/activities/common.ts index 6e8b463267..93afbe5605 100644 --- a/services/apps/script_executor_worker/src/activities/common.ts +++ b/services/apps/script_executor_worker/src/activities/common.ts @@ -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, @@ -27,6 +28,57 @@ export async function mergeMembers( } } +export async function mergeMembersIfAllowed( + primaryMemberId: string, + secondaryMemberId: string, +): Promise { + 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), + ) + + 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 | IUnmergePreviewResult, diff --git a/services/apps/script_executor_worker/src/activities/merge-members-with-similar-identities/index.ts b/services/apps/script_executor_worker/src/activities/merge-members-with-similar-identities/index.ts index e836e2269d..a0ea7f2843 100644 --- a/services/apps/script_executor_worker/src/activities/merge-members-with-similar-identities/index.ts +++ b/services/apps/script_executor_worker/src/activities/merge-members-with-similar-identities/index.ts @@ -5,13 +5,18 @@ import { svc } from '../../main' export async function findMembersWithSameVerifiedEmailsInDifferentPlatforms( limit: number, - afterHash?: number, + afterHighMemberId?: string, + afterLowMemberId?: string, ): Promise { 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) } diff --git a/services/apps/script_executor_worker/src/main.ts b/services/apps/script_executor_worker/src/main.ts index 645f060c34..147f39ba84 100644 --- a/services/apps/script_executor_worker/src/main.ts +++ b/services/apps/script_executor_worker/src/main.ts @@ -7,6 +7,7 @@ import { scheduleOrganizationSegmentAggCleanup, scheduleOrganizationsCleanup, } from './schedules/scheduleCleanup' +import { scheduleMergeMembersWithSameVerifiedEmails } from './schedules/scheduleMemberDeduplication' const config: Config = { envvars: [ @@ -47,6 +48,7 @@ setImmediate(async () => { await scheduleOrganizationsCleanup() await scheduleMemberSegmentsAggCleanup() await scheduleOrganizationSegmentAggCleanup() + await scheduleMergeMembersWithSameVerifiedEmails() await svc.start() }) diff --git a/services/apps/script_executor_worker/src/schedules/scheduleMemberDeduplication.ts b/services/apps/script_executor_worker/src/schedules/scheduleMemberDeduplication.ts new file mode 100644 index 0000000000..d0db325d63 --- /dev/null +++ b/services/apps/script_executor_worker/src/schedules/scheduleMemberDeduplication.ts @@ -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, + 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) + } + } +} diff --git a/services/apps/script_executor_worker/src/types.ts b/services/apps/script_executor_worker/src/types.ts index 428d712ab7..04d460455a 100644 --- a/services/apps/script_executor_worker/src/types.ts +++ b/services/apps/script_executor_worker/src/types.ts @@ -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 { diff --git a/services/apps/script_executor_worker/src/workflows/findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatforms.ts b/services/apps/script_executor_worker/src/workflows/findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatforms.ts index fff1a3f60a..da1c66f8d9 100644 --- a/services/apps/script_executor_worker/src/workflows/findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatforms.ts +++ b/services/apps/script_executor_worker/src/workflows/findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatforms.ts @@ -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) { @@ -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({ - afterHash: mergeableMemberCouples[mergeableMemberCouples.length - 1]?.hash, + afterHighMemberId, + afterLowMemberId, + dryRun: args.dryRun, }) } diff --git a/services/libs/common_services/src/services/common.member.service.ts b/services/libs/common_services/src/services/common.member.service.ts index 0c0bf069a5..3107de55cb 100644 --- a/services/libs/common_services/src/services/common.member.service.ts +++ b/services/libs/common_services/src/services/common.member.service.ts @@ -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, @@ -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) diff --git a/services/libs/data-access-layer/src/member_merge/index.ts b/services/libs/data-access-layer/src/member_merge/index.ts index 18b0c975f5..0ac18da99b 100644 --- a/services/libs/data-access-layer/src/member_merge/index.ts +++ b/services/libs/data-access-layer/src/member_merge/index.ts @@ -50,6 +50,33 @@ export async function insertMemberNoMerge( ) } +export async function moveMemberNoMerge( + qx: QueryExecutor, + fromMemberId: string, + toMemberId: string, +): Promise { + 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[], diff --git a/services/libs/data-access-layer/src/members/base.ts b/services/libs/data-access-layer/src/members/base.ts index abebc58303..0735bc8f01 100644 --- a/services/libs/data-access-layer/src/members/base.ts +++ b/services/libs/data-access-layer/src/members/base.ts @@ -716,6 +716,17 @@ export async function findMemberById( return queryTableById(qx, 'members', Object.values(MemberField), memberId, fields) } +export async function findExistingMemberIds( + qx: QueryExecutor, + memberIds: string[], +): Promise { + 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, diff --git a/services/libs/data-access-layer/src/old/apps/script_executor_worker/member.repo.ts b/services/libs/data-access-layer/src/old/apps/script_executor_worker/member.repo.ts index b9fb5dbffa..b9c252b8dc 100644 --- a/services/libs/data-access-layer/src/old/apps/script_executor_worker/member.repo.ts +++ b/services/libs/data-access-layer/src/old/apps/script_executor_worker/member.repo.ts @@ -16,13 +16,15 @@ class MemberRepository { async findMembersWithSameVerifiedEmailsInDifferentPlatforms( limit = 50, - afterHash: number = undefined, + afterHighMemberId: string = undefined, + afterLowMemberId: string = undefined, ): Promise { 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 @@ -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 + order by + Greatest(a."memberId"::text, b."memberId"::text) desc, + Least(a."memberId"::text, b."memberId"::text) desc limit $(limit); `, { - afterHash, + afterHighMemberId, + afterLowMemberId, limit, }, )