From 69d619c26a5c79561edd222e6471fc5bafc8aa88 Mon Sep 17 00:00:00 2001 From: Ramanathan Date: Mon, 17 Aug 2026 18:44:16 +0530 Subject: [PATCH 1/6] fix(script-executor): schedule member email deduplication sweep findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatforms is registered in workflows.ts but nothing ever triggers it. main.ts registers only the four cleanup schedules, and services/cronjobs holds only archived_repositories, so the sweep runs only when someone starts it by hand and duplicate members accumulate in between. Register it as a recurring schedule alongside the existing ones. The workflow already pages by hash and uses continueAsNew, so it is built for unattended repeated execution. Matching semantics are unchanged: the detection query still requires verified identities on both sides, so this adds no new trust in self-asserted git commit emails. Overlap policy is SKIP rather than BUFFER_ONE because a full sweep can outlast the interval, and buffered sweeps would stack up. Refs: https://github.com/linuxfoundation/crowd.dev/issues/4484 Co-Authored-By: Claude Opus 5 (1M context) Signed-off-by: Ramanathan --- .../apps/script_executor_worker/src/main.ts | 2 + .../schedules/scheduleMemberDeduplication.ts | 44 +++++++++++++++++++ 2 files changed, 46 insertions(+) create mode 100644 services/apps/script_executor_worker/src/schedules/scheduleMemberDeduplication.ts 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..ae00ae8d09 --- /dev/null +++ b/services/apps/script_executor_worker/src/schedules/scheduleMemberDeduplication.ts @@ -0,0 +1,44 @@ +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: { + // Run every Sunday at 03:00, off-peak and clear of the Wednesday cleanup schedules + cronExpressions: ['0 3 * * 0'], + }, + policies: { + // The workflow walks the whole memberIdentities table via continueAsNew, so a single + // sweep can outlast the interval. Skip an overdue run instead of buffering it, so + // sweeps never stack up on top of each other. + overlap: ScheduleOverlapPolicy.SKIP, + catchupWindow: '1 minute', + }, + action: { + type: 'startWorkflow', + workflowType: findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatforms, + taskQueue: 'script-executor', + retry: { + initialInterval: '15 seconds', + backoffCoefficient: 2, + maximumAttempts: 3, + }, + // Start from the beginning of the hash ordering; the workflow pages itself from there. + 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) + } + } +} From 94b5d95b7274f0adecf600b4c273956c03e4e4af Mon Sep 17 00:00:00 2001 From: Ramanathan Date: Mon, 17 Aug 2026 20:07:10 +0530 Subject: [PATCH 2/6] Potential fix for pull request finding removing comments Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> Signed-off-by: Ramanathan --- .../src/schedules/scheduleMemberDeduplication.ts | 1 - 1 file changed, 1 deletion(-) diff --git a/services/apps/script_executor_worker/src/schedules/scheduleMemberDeduplication.ts b/services/apps/script_executor_worker/src/schedules/scheduleMemberDeduplication.ts index ae00ae8d09..ac66c249fc 100644 --- a/services/apps/script_executor_worker/src/schedules/scheduleMemberDeduplication.ts +++ b/services/apps/script_executor_worker/src/schedules/scheduleMemberDeduplication.ts @@ -8,7 +8,6 @@ export const scheduleMergeMembersWithSameVerifiedEmails = async () => { await svc.temporal.schedule.create({ scheduleId: 'mergeMembersWithSameVerifiedEmails', spec: { - // Run every Sunday at 03:00, off-peak and clear of the Wednesday cleanup schedules cronExpressions: ['0 3 * * 0'], }, policies: { From 6c3a4482d773344e71f8659fb5aa1e99cadb7989 Mon Sep 17 00:00:00 2001 From: Ramanathan Date: Tue, 18 Aug 2026 15:39:48 +0530 Subject: [PATCH 3/6] fix(script-executor): guard member email deduplication sweep Respect memberNoMerge, isolate per-couple failures, add dry run. Signed-off-by: Ramanathan --- .../script_executor_worker/src/activities.ts | 2 ++ .../src/activities/common.ts | 34 +++++++++++++++++++ .../schedules/scheduleMemberDeduplication.ts | 4 --- .../apps/script_executor_worker/src/types.ts | 1 + ...hSameVerifiedEmailsInDifferentPlatforms.ts | 19 ++++++++--- 5 files changed, 52 insertions(+), 8 deletions(-) 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..2baa57e0ee 100644 --- a/services/apps/script_executor_worker/src/activities/common.ts +++ b/services/apps/script_executor_worker/src/activities/common.ts @@ -2,6 +2,7 @@ import axios from 'axios' import { CommonMemberService, signalMemberUpdate } from '@crowd/common_services' import { pgpQx } from '@crowd/data-access-layer' +import { getMemberNoMerge } from '@crowd/data-access-layer/src/member_merge' import { IMemberIdentity, IMemberUnmergeBackup, @@ -27,6 +28,39 @@ export async function mergeMembers( } } +export async function mergeMembersIfAllowed( + primaryMemberId: string, + secondaryMemberId: string, +): Promise { + const qx = pgpQx(svc.postgres.writer.connection()) + + 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) { + 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/schedules/scheduleMemberDeduplication.ts b/services/apps/script_executor_worker/src/schedules/scheduleMemberDeduplication.ts index ac66c249fc..d0db325d63 100644 --- a/services/apps/script_executor_worker/src/schedules/scheduleMemberDeduplication.ts +++ b/services/apps/script_executor_worker/src/schedules/scheduleMemberDeduplication.ts @@ -11,9 +11,6 @@ export const scheduleMergeMembersWithSameVerifiedEmails = async () => { cronExpressions: ['0 3 * * 0'], }, policies: { - // The workflow walks the whole memberIdentities table via continueAsNew, so a single - // sweep can outlast the interval. Skip an overdue run instead of buffering it, so - // sweeps never stack up on top of each other. overlap: ScheduleOverlapPolicy.SKIP, catchupWindow: '1 minute', }, @@ -26,7 +23,6 @@ export const scheduleMergeMembersWithSameVerifiedEmails = async () => { backoffCoefficient: 2, maximumAttempts: 3, }, - // Start from the beginning of the hash ordering; the workflow pages itself from there. args: [{}], }, }) diff --git a/services/apps/script_executor_worker/src/types.ts b/services/apps/script_executor_worker/src/types.ts index 428d712ab7..7e04edba5f 100644 --- a/services/apps/script_executor_worker/src/types.ts +++ b/services/apps/script_executor_worker/src/types.ts @@ -2,6 +2,7 @@ import { IActivityRelationDuplicateGroup } from '@crowd/data-access-layer' export interface IFindAndMergeMembersWithSameVerifiedEmailsInDifferentPlatformsArgs { afterHash?: number + 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..1061b77df1 100644 --- a/services/apps/script_executor_worker/src/workflows/findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatforms.ts +++ b/services/apps/script_executor_worker/src/workflows/findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatforms.ts @@ -31,13 +31,24 @@ 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}!`) + + try { + await common.mergeMembersIfAllowed(couple.primaryMemberId, couple.secondaryMemberId) + } catch (err) { + console.log(`Failed to merge ${coupleDescription}, skipping couple! Error: ${err?.message}`) + } } await continueAsNew({ afterHash: mergeableMemberCouples[mergeableMemberCouples.length - 1]?.hash, + dryRun: args.dryRun, }) } From d757758b4289f95b4d1dd366e1e4387f3170549d Mon Sep 17 00:00:00 2001 From: Ramanathan Date: Tue, 18 Aug 2026 15:51:04 +0530 Subject: [PATCH 4/6] fix(script-executor): carry no-merge edges when absorbing a member Direct-pair checks alone let a blocked pair rejoin transitively through a third member. Signed-off-by: Ramanathan --- .../src/activities/common.ts | 4 ++- .../src/member_merge/index.ts | 27 +++++++++++++++++++ 2 files changed, 30 insertions(+), 1 deletion(-) diff --git a/services/apps/script_executor_worker/src/activities/common.ts b/services/apps/script_executor_worker/src/activities/common.ts index 2baa57e0ee..d0e393449b 100644 --- a/services/apps/script_executor_worker/src/activities/common.ts +++ b/services/apps/script_executor_worker/src/activities/common.ts @@ -2,7 +2,7 @@ import axios from 'axios' import { CommonMemberService, signalMemberUpdate } from '@crowd/common_services' import { pgpQx } from '@crowd/data-access-layer' -import { getMemberNoMerge } from '@crowd/data-access-layer/src/member_merge' +import { getMemberNoMerge, moveMemberNoMerge } from '@crowd/data-access-layer/src/member_merge' import { IMemberIdentity, IMemberUnmergeBackup, @@ -49,6 +49,8 @@ export async function mergeMembersIfAllowed( return false } + await moveMemberNoMerge(qx, secondaryMemberId, primaryMemberId) + const memberService = new CommonMemberService(qx, svc.temporal, svc.log) try { 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[], From 70eed0233b764cb901616454ab821f23be50ffa5 Mon Sep 17 00:00:00 2001 From: Ramanathan Date: Tue, 18 Aug 2026 15:55:34 +0530 Subject: [PATCH 5/6] Potential fix for pull request finding updated Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> Signed-off-by: Ramanathan --- .../script_executor_worker/src/activities/common.ts | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/services/apps/script_executor_worker/src/activities/common.ts b/services/apps/script_executor_worker/src/activities/common.ts index d0e393449b..4864f8b069 100644 --- a/services/apps/script_executor_worker/src/activities/common.ts +++ b/services/apps/script_executor_worker/src/activities/common.ts @@ -49,12 +49,12 @@ export async function mergeMembersIfAllowed( return false } - await moveMemberNoMerge(qx, secondaryMemberId, primaryMemberId) - - const memberService = new CommonMemberService(qx, svc.temporal, svc.log) - try { - await memberService.merge(primaryMemberId, secondaryMemberId) + await qx.tx(async (txQx) => { + await moveMemberNoMerge(txQx, secondaryMemberId, primaryMemberId) + const memberService = new CommonMemberService(txQx, svc.temporal, svc.log) + await memberService.merge(primaryMemberId, secondaryMemberId) + }) } catch (error) { svc.log.error({ err: error, primaryMemberId, secondaryMemberId }, 'Failed to merge members') throw error From d6dc875f3d0e2e6a1cbf866af99c4086c80b5469 Mon Sep 17 00:00:00 2001 From: Ramanathan Date: Tue, 18 Aug 2026 16:05:08 +0530 Subject: [PATCH 6/6] fix(script-executor): make no-merge propagation atomic with the merge Classify expected skips in the activity so real failures still fail the sweep. Signed-off-by: Ramanathan --- .../src/activities/common.ts | 30 ++++++++++++++----- ...hSameVerifiedEmailsInDifferentPlatforms.ts | 6 +--- .../src/services/common.member.service.ts | 4 ++- .../data-access-layer/src/members/base.ts | 11 +++++++ 4 files changed, 38 insertions(+), 13 deletions(-) diff --git a/services/apps/script_executor_worker/src/activities/common.ts b/services/apps/script_executor_worker/src/activities/common.ts index 4864f8b069..93afbe5605 100644 --- a/services/apps/script_executor_worker/src/activities/common.ts +++ b/services/apps/script_executor_worker/src/activities/common.ts @@ -1,8 +1,8 @@ import axios from 'axios' import { CommonMemberService, signalMemberUpdate } from '@crowd/common_services' -import { pgpQx } from '@crowd/data-access-layer' -import { getMemberNoMerge, moveMemberNoMerge } from '@crowd/data-access-layer/src/member_merge' +import { findExistingMemberIds, pgpQx } from '@crowd/data-access-layer' +import { getMemberNoMerge } from '@crowd/data-access-layer/src/member_merge' import { IMemberIdentity, IMemberUnmergeBackup, @@ -34,6 +34,16 @@ export async function mergeMembersIfAllowed( ): 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) => @@ -49,13 +59,19 @@ export async function mergeMembersIfAllowed( return false } + const memberService = new CommonMemberService(qx, svc.temporal, svc.log) + try { - await qx.tx(async (txQx) => { - await moveMemberNoMerge(txQx, secondaryMemberId, primaryMemberId) - const memberService = new CommonMemberService(txQx, svc.temporal, svc.log) - await memberService.merge(primaryMemberId, secondaryMemberId) - }) + 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 } diff --git a/services/apps/script_executor_worker/src/workflows/findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatforms.ts b/services/apps/script_executor_worker/src/workflows/findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatforms.ts index 1061b77df1..e5f7b05141 100644 --- a/services/apps/script_executor_worker/src/workflows/findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatforms.ts +++ b/services/apps/script_executor_worker/src/workflows/findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatforms.ts @@ -40,11 +40,7 @@ export async function findAndMergeMembersWithSameVerifiedEmailsInDifferentPlatfo console.log(`Merging ${coupleDescription}!`) - try { - await common.mergeMembersIfAllowed(couple.primaryMemberId, couple.secondaryMemberId) - } catch (err) { - console.log(`Failed to merge ${coupleDescription}, skipping couple! Error: ${err?.message}`) - } + await common.mergeMembersIfAllowed(couple.primaryMemberId, couple.secondaryMemberId) } await continueAsNew({ 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/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,