Skip to content

Commit e021298

Browse files
authored
fix: keep memberToMerge and memberToMergeRaw in sync (CM-1339) (#4384)
Signed-off-by: Yeganathan S <63534555+skwowet@users.noreply.github.com>
1 parent 177c7d2 commit e021298

8 files changed

Lines changed: 52 additions & 176 deletions

File tree

backend/src/database/repositories/memberRepository.ts

Lines changed: 4 additions & 70 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import lodash, { chunk, uniq } from 'lodash'
1+
import lodash, { uniq } from 'lodash'
22
import Sequelize, { QueryTypes } from 'sequelize'
33

44
import {
@@ -29,7 +29,7 @@ import {
2929
} from '@crowd/data-access-layer'
3030
import { findManyLfxMemberships } from '@crowd/data-access-layer/src/lfx_memberships'
3131
import { findMaintainerRoles } from '@crowd/data-access-layer/src/maintainers'
32-
import { addMemberNoMerge, removeMemberToMerge } from '@crowd/data-access-layer/src/member_merge'
32+
import { insertMemberNoMerge, removeMemberToMerge } from '@crowd/data-access-layer/src/member_merge'
3333
import {
3434
deleteMemberSegmentAffiliations,
3535
findMemberAffiliations,
@@ -86,7 +86,7 @@ import MemberAttributeSettingsRepository from './memberAttributeSettingsReposito
8686
import SegmentRepository from './segmentRepository'
8787
import SequelizeRepository from './sequelizeRepository'
8888
import TenantRepository from './tenantRepository'
89-
import { IMemberMergeSuggestion, mapUsernameToIdentities } from './types/memberTypes'
89+
import { mapUsernameToIdentities } from './types/memberTypes'
9090

9191
const { Op } = Sequelize
9292

@@ -581,72 +581,6 @@ class MemberRepository {
581581
}
582582
}
583583

584-
static async addToMerge(
585-
suggestions: IMemberMergeSuggestion[],
586-
options: IRepositoryOptions,
587-
): Promise<void> {
588-
const transaction = SequelizeRepository.getTransaction(options)
589-
const seq = SequelizeRepository.getSequelize(options)
590-
591-
// Remove possible duplicates
592-
suggestions = lodash.uniqWith(suggestions, (a, b) =>
593-
lodash.isEqual(lodash.sortBy(a.members), lodash.sortBy(b.members)),
594-
)
595-
596-
// Process suggestions in chunks of 100 or less
597-
const suggestionChunks = chunk(suggestions, 100)
598-
599-
const insertValues = (
600-
memberId: string,
601-
toMergeId: string,
602-
similarity: number | null,
603-
index: number,
604-
) => {
605-
const idPlaceholder = (key: string) => `${key}${index}`
606-
return {
607-
query: `(:${idPlaceholder('memberId')}, :${idPlaceholder('toMergeId')}, :${idPlaceholder(
608-
'similarity',
609-
)}, NOW(), NOW())`,
610-
replacements: {
611-
[idPlaceholder('memberId')]: memberId,
612-
[idPlaceholder('toMergeId')]: toMergeId,
613-
[idPlaceholder('similarity')]: similarity === null ? null : similarity,
614-
},
615-
}
616-
}
617-
618-
for (const suggestionChunk of suggestionChunks) {
619-
const placeholders: string[] = []
620-
let replacements: Record<string, unknown> = {}
621-
622-
suggestionChunk.forEach((suggestion, index) => {
623-
const { query, replacements: chunkReplacements } = insertValues(
624-
suggestion.members[0],
625-
suggestion.members[1],
626-
suggestion.similarity,
627-
index,
628-
)
629-
placeholders.push(query)
630-
replacements = { ...replacements, ...chunkReplacements }
631-
})
632-
633-
const query = `
634-
INSERT INTO "memberToMerge" ("memberId", "toMergeId", "similarity", "createdAt", "updatedAt")
635-
VALUES ${placeholders.join(', ')} on conflict do nothing;
636-
`
637-
try {
638-
await seq.query(query, {
639-
replacements,
640-
type: QueryTypes.INSERT,
641-
transaction,
642-
})
643-
} catch (error) {
644-
options.log.error('error adding members to merge', error)
645-
throw error
646-
}
647-
}
648-
}
649-
650584
static async removeToMerge(id, toMergeId, options: IRepositoryOptions) {
651585
const qx = SequelizeRepository.getQueryExecutor(options)
652586

@@ -656,7 +590,7 @@ class MemberRepository {
656590
static async addNoMerge(id, toMergeId, options: IRepositoryOptions) {
657591
const qx = SequelizeRepository.getQueryExecutor(options)
658592

659-
await addMemberNoMerge(qx, id, toMergeId)
593+
await insertMemberNoMerge(qx, id, toMergeId)
660594
}
661595

662596
static async memberExists(

backend/src/services/memberService.ts

Lines changed: 2 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -49,7 +49,6 @@ import { MergeActionsRepository } from '../database/repositories/mergeActionsRep
4949
import SequelizeRepository from '../database/repositories/sequelizeRepository'
5050
import {
5151
BasicMemberIdentity,
52-
IMemberMergeSuggestion,
5352
mapUsernameToIdentities,
5453
} from '../database/repositories/types/memberTypes'
5554
import telemetryTrack from '../segment/telemetryTrack'
@@ -729,32 +728,6 @@ export default class MemberService extends LoggerBase {
729728
return data
730729
}
731730

732-
/**
733-
* Given two members, add them to the toMerge fields of each other.
734-
* It will also update the tenant's toMerge list, removing any entry that contains
735-
* the pair.
736-
* @returns Success/Error message
737-
*/
738-
async addToMerge(suggestions: IMemberMergeSuggestion[]) {
739-
const transaction = await SequelizeRepository.createTransaction(this.options)
740-
try {
741-
const searchSyncService = new SearchSyncService(this.options)
742-
743-
await MemberRepository.addToMerge(suggestions, { ...this.options, transaction })
744-
await SequelizeRepository.commitTransaction(transaction)
745-
746-
for (const suggestion of suggestions) {
747-
await searchSyncService.triggerMemberSync(suggestion.members[0])
748-
await searchSyncService.triggerMemberSync(suggestion.members[1])
749-
}
750-
return { status: 200 }
751-
} catch (error) {
752-
await SequelizeRepository.rollbackTransaction(transaction)
753-
this.log.error(error, 'Error while adding members to merge')
754-
throw error
755-
}
756-
}
757-
758731
/**
759732
* Given two members, add them to the noMerge fields of each other.
760733
* @param memberOneId ID of the first member
@@ -768,8 +741,9 @@ export default class MemberService extends LoggerBase {
768741
try {
769742
await MemberRepository.addNoMerge(memberOneId, memberTwoId, txOptions)
770743
await MemberRepository.addNoMerge(memberTwoId, memberOneId, txOptions)
744+
745+
// Removes from either order of the pair
771746
await MemberRepository.removeToMerge(memberOneId, memberTwoId, txOptions)
772-
await MemberRepository.removeToMerge(memberTwoId, memberOneId, txOptions)
773747

774748
await SequelizeRepository.commitTransaction(transaction)
775749
} catch (error) {

services/apps/merge_suggestions_worker/src/activities/memberMergeSuggestions.ts

Lines changed: 11 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22
import uniqBy from 'lodash.uniqby'
33

44
import { parseGitHubNoreplyEmail, parseGitLabNoreplyEmail } from '@crowd/common'
5-
import { addMemberNoMerge } from '@crowd/data-access-layer/src/member_merge'
5+
import { insertMemberNoMerge, removeMemberToMerge } from '@crowd/data-access-layer/src/member_merge'
66
import { MemberField, queryMembers } from '@crowd/data-access-layer/src/members'
77
import MemberMergeSuggestionsRepository from '@crowd/data-access-layer/src/old/apps/merge_suggestions_worker/memberMergeSuggestions.repo'
88
import { pgpQx } from '@crowd/data-access-layer/src/queryExecutor'
@@ -498,15 +498,14 @@ export async function getRawMemberMergeSuggestions(
498498
return memberMergeSuggestionsRepo.getRawMemberSuggestions(similarityFilter, limit)
499499
}
500500

501-
export async function removeMemberMergeSuggestion(
502-
suggestion: string[],
503-
table: MemberMergeSuggestionTable,
504-
): Promise<void> {
505-
const memberMergeSuggestionsRepo = new MemberMergeSuggestionsRepository(
506-
svc.postgres.writer.connection(),
507-
svc.log,
508-
)
509-
await memberMergeSuggestionsRepo.removeMemberMergeSuggestion(suggestion, table)
501+
export async function removeMemberMergeSuggestion(suggestion: string[]): Promise<void> {
502+
if (suggestion.length !== 2) {
503+
svc.log.debug(`Suggestions array must have two ids!`)
504+
return
505+
}
506+
507+
const qx = pgpQx(svc.postgres.writer.connection())
508+
await removeMemberToMerge(qx, suggestion[0], suggestion[1])
510509
}
511510

512511
export async function addMemberSuggestionToNoMerge(suggestion: string[]): Promise<void> {
@@ -516,5 +515,6 @@ export async function addMemberSuggestionToNoMerge(suggestion: string[]): Promis
516515
}
517516
const qx = pgpQx(svc.postgres.writer.connection())
518517

519-
await addMemberNoMerge(qx, suggestion[0], suggestion[1])
518+
await insertMemberNoMerge(qx, suggestion[0], suggestion[1])
519+
await insertMemberNoMerge(qx, suggestion[1], suggestion[0])
520520
}

services/apps/merge_suggestions_worker/src/workflows/mergeMembersWithLLM.ts

Lines changed: 3 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import { continueAsNew, proxyActivities } from '@temporalio/workflow'
22

3-
import { LLMSuggestionVerdictType, MemberMergeSuggestionTable } from '@crowd/types'
3+
import { LLMSuggestionVerdictType } from '@crowd/types'
44

55
import type * as activities from '../activities'
66
import { ILLMResult, IProcessMergeMemberSuggestionsWithLLM } from '../types'
@@ -78,11 +78,7 @@ export async function mergeMembersWithLLM(
7878

7979
if (members.length !== 2) {
8080
console.log(`Failed getting members data in suggestion. Skipping suggestion: ${suggestion}`)
81-
await removeMemberMergeSuggestion(suggestion, MemberMergeSuggestionTable.MEMBER_TO_MERGE_RAW)
82-
await removeMemberMergeSuggestion(
83-
suggestion,
84-
MemberMergeSuggestionTable.MEMBER_TO_MERGE_FILTERED,
85-
)
81+
await removeMemberMergeSuggestion(suggestion)
8682
continue
8783
}
8884

@@ -115,11 +111,7 @@ export async function mergeMembersWithLLM(
115111
console.log(
116112
`LLM doesn't think these members are the same. Removing from suggestions and adding to no merge: ${suggestion[0]} and ${suggestion[1]}!`,
117113
)
118-
await removeMemberMergeSuggestion(
119-
suggestion,
120-
MemberMergeSuggestionTable.MEMBER_TO_MERGE_FILTERED,
121-
)
122-
await removeMemberMergeSuggestion(suggestion, MemberMergeSuggestionTable.MEMBER_TO_MERGE_RAW)
114+
await removeMemberMergeSuggestion(suggestion)
123115
await addMemberSuggestionToNoMerge(suggestion)
124116
}
125117
}

services/libs/common_services/src/services/member/unmerge.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@ import {
2121
removeMemberRole,
2222
updateMember,
2323
} from '@crowd/data-access-layer'
24-
import { addMemberNoMerge } from '@crowd/data-access-layer/src/member_merge'
24+
import { insertMemberNoMerge } from '@crowd/data-access-layer/src/member_merge'
2525
import {
2626
deleteMemberSegmentAffiliations,
2727
findMemberAffiliations,
@@ -625,7 +625,7 @@ export async function unmergeMember(
625625
}
626626

627627
// Add primary and secondary to no merge so they don't get suggested again
628-
await addMemberNoMerge(tx, memberId, secondaryId)
628+
await insertMemberNoMerge(tx, memberId, secondaryId)
629629

630630
await setMergeAction(tx, MergeActionType.MEMBER, memberId, secondaryId, {
631631
step: MergeActionStep.UNMERGE_SYNC_DONE,

services/libs/data-access-layer/src/member_merge/index.ts

Lines changed: 30 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -5,37 +5,46 @@ export async function removeMemberToMerge(
55
memberId: string,
66
toMergeId: string,
77
): Promise<void> {
8-
await qx.result(
9-
`
10-
DELETE FROM "memberToMerge"
11-
WHERE "memberId" = $(memberId)
12-
AND "toMergeId" = $(toMergeId)
13-
`,
14-
{
15-
memberId,
16-
toMergeId,
17-
},
18-
)
8+
const replacements = { memberId, toMergeId }
9+
10+
const whereClause = `
11+
WHERE
12+
("memberId" = $(memberId) AND "toMergeId" = $(toMergeId))
13+
OR
14+
("memberId" = $(toMergeId) AND "toMergeId" = $(memberId))
15+
`
16+
17+
await qx.tx(async (tx) => {
18+
await tx.result(
19+
`
20+
DELETE FROM "memberToMerge"
21+
${whereClause}
22+
`,
23+
replacements,
24+
)
25+
26+
await tx.result(
27+
`
28+
DELETE FROM "memberToMergeRaw"
29+
${whereClause}
30+
`,
31+
replacements,
32+
)
33+
})
1934
}
2035

21-
export async function addMemberNoMerge(
36+
export async function insertMemberNoMerge(
2237
qx: QueryExecutor,
2338
memberId: string,
2439
noMergeId: string,
2540
): Promise<void> {
26-
const currentTime = new Date()
2741
await qx.result(
2842
`
2943
INSERT INTO "memberNoMerge" ("memberId", "noMergeId", "createdAt", "updatedAt")
30-
VALUES ($(memberId), $(noMergeId), $(createdAt), $(updatedAt))
31-
on conflict ("memberId", "noMergeId") do nothing
44+
VALUES ($(memberId), $(noMergeId), NOW(), NOW())
45+
ON CONFLICT ("memberId", "noMergeId") DO NOTHING
3246
`,
33-
{
34-
memberId,
35-
noMergeId,
36-
createdAt: currentTime,
37-
updatedAt: currentTime,
38-
},
47+
{ memberId, noMergeId },
3948
)
4049
}
4150

services/libs/data-access-layer/src/old/apps/members_enrichment_worker/index.ts

Lines changed: 0 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -365,14 +365,6 @@ export async function findExistingMember(
365365
return results.map((r) => r.memberId)
366366
}
367367

368-
export async function addMemberToMerge(tx: DbTransaction, memberId: string, toMergeId: string) {
369-
await tx.query(
370-
`INSERT INTO "memberToMerge" ("memberId", "toMergeId", similarity)
371-
VALUES ($1, $2, $3);"`,
372-
[memberId, toMergeId, 0.9],
373-
)
374-
}
375-
376368
export async function findOrganizationIdentities(
377369
tx: DbTransaction,
378370
organizationId: string,

services/libs/data-access-layer/src/old/apps/merge_suggestions_worker/memberMergeSuggestions.repo.ts

Lines changed: 0 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -325,31 +325,6 @@ class MemberMergeSuggestionsRepository {
325325

326326
return results.map((r) => [r.memberId, r.toMergeId])
327327
}
328-
329-
async removeMemberMergeSuggestion(
330-
suggestion: string[],
331-
table: MemberMergeSuggestionTable,
332-
): Promise<void> {
333-
const query = `
334-
delete from "${table}"
335-
where
336-
("memberId" = $(memberId) and "toMergeId" = $(toMergeId))
337-
or
338-
("memberId" = $(toMergeId) and "toMergeId" = $(memberId))
339-
`
340-
341-
const replacements = {
342-
memberId: suggestion[0],
343-
toMergeId: suggestion[1],
344-
}
345-
346-
try {
347-
await this.connection.none(query, replacements)
348-
} catch (error) {
349-
this.log.error(`Error removing member suggestions from ${table}`, error)
350-
throw error
351-
}
352-
}
353328
}
354329

355330
export default MemberMergeSuggestionsRepository

0 commit comments

Comments
 (0)