diff --git a/backend/src/database/repositories/memberRepository.ts b/backend/src/database/repositories/memberRepository.ts index 7e6f69a5a1..a9f86dc023 100644 --- a/backend/src/database/repositories/memberRepository.ts +++ b/backend/src/database/repositories/memberRepository.ts @@ -1,4 +1,4 @@ -import lodash, { chunk, uniq } from 'lodash' +import lodash, { uniq } from 'lodash' import Sequelize, { QueryTypes } from 'sequelize' import { @@ -29,7 +29,7 @@ import { } from '@crowd/data-access-layer' import { findManyLfxMemberships } from '@crowd/data-access-layer/src/lfx_memberships' import { findMaintainerRoles } from '@crowd/data-access-layer/src/maintainers' -import { addMemberNoMerge, removeMemberToMerge } from '@crowd/data-access-layer/src/member_merge' +import { insertMemberNoMerge, removeMemberToMerge } from '@crowd/data-access-layer/src/member_merge' import { deleteMemberSegmentAffiliations, findMemberAffiliations, @@ -86,7 +86,7 @@ import MemberAttributeSettingsRepository from './memberAttributeSettingsReposito import SegmentRepository from './segmentRepository' import SequelizeRepository from './sequelizeRepository' import TenantRepository from './tenantRepository' -import { IMemberMergeSuggestion, mapUsernameToIdentities } from './types/memberTypes' +import { mapUsernameToIdentities } from './types/memberTypes' const { Op } = Sequelize @@ -581,72 +581,6 @@ class MemberRepository { } } - static async addToMerge( - suggestions: IMemberMergeSuggestion[], - options: IRepositoryOptions, - ): Promise { - const transaction = SequelizeRepository.getTransaction(options) - const seq = SequelizeRepository.getSequelize(options) - - // Remove possible duplicates - suggestions = lodash.uniqWith(suggestions, (a, b) => - lodash.isEqual(lodash.sortBy(a.members), lodash.sortBy(b.members)), - ) - - // Process suggestions in chunks of 100 or less - const suggestionChunks = chunk(suggestions, 100) - - const insertValues = ( - memberId: string, - toMergeId: string, - similarity: number | null, - index: number, - ) => { - const idPlaceholder = (key: string) => `${key}${index}` - return { - query: `(:${idPlaceholder('memberId')}, :${idPlaceholder('toMergeId')}, :${idPlaceholder( - 'similarity', - )}, NOW(), NOW())`, - replacements: { - [idPlaceholder('memberId')]: memberId, - [idPlaceholder('toMergeId')]: toMergeId, - [idPlaceholder('similarity')]: similarity === null ? null : similarity, - }, - } - } - - for (const suggestionChunk of suggestionChunks) { - const placeholders: string[] = [] - let replacements: Record = {} - - suggestionChunk.forEach((suggestion, index) => { - const { query, replacements: chunkReplacements } = insertValues( - suggestion.members[0], - suggestion.members[1], - suggestion.similarity, - index, - ) - placeholders.push(query) - replacements = { ...replacements, ...chunkReplacements } - }) - - const query = ` - INSERT INTO "memberToMerge" ("memberId", "toMergeId", "similarity", "createdAt", "updatedAt") - VALUES ${placeholders.join(', ')} on conflict do nothing; - ` - try { - await seq.query(query, { - replacements, - type: QueryTypes.INSERT, - transaction, - }) - } catch (error) { - options.log.error('error adding members to merge', error) - throw error - } - } - } - static async removeToMerge(id, toMergeId, options: IRepositoryOptions) { const qx = SequelizeRepository.getQueryExecutor(options) @@ -656,7 +590,7 @@ class MemberRepository { static async addNoMerge(id, toMergeId, options: IRepositoryOptions) { const qx = SequelizeRepository.getQueryExecutor(options) - await addMemberNoMerge(qx, id, toMergeId) + await insertMemberNoMerge(qx, id, toMergeId) } static async memberExists( diff --git a/backend/src/services/memberService.ts b/backend/src/services/memberService.ts index 052e4f3b33..a9b1fdd164 100644 --- a/backend/src/services/memberService.ts +++ b/backend/src/services/memberService.ts @@ -49,7 +49,6 @@ import { MergeActionsRepository } from '../database/repositories/mergeActionsRep import SequelizeRepository from '../database/repositories/sequelizeRepository' import { BasicMemberIdentity, - IMemberMergeSuggestion, mapUsernameToIdentities, } from '../database/repositories/types/memberTypes' import telemetryTrack from '../segment/telemetryTrack' @@ -729,32 +728,6 @@ export default class MemberService extends LoggerBase { return data } - /** - * Given two members, add them to the toMerge fields of each other. - * It will also update the tenant's toMerge list, removing any entry that contains - * the pair. - * @returns Success/Error message - */ - async addToMerge(suggestions: IMemberMergeSuggestion[]) { - const transaction = await SequelizeRepository.createTransaction(this.options) - try { - const searchSyncService = new SearchSyncService(this.options) - - await MemberRepository.addToMerge(suggestions, { ...this.options, transaction }) - await SequelizeRepository.commitTransaction(transaction) - - for (const suggestion of suggestions) { - await searchSyncService.triggerMemberSync(suggestion.members[0]) - await searchSyncService.triggerMemberSync(suggestion.members[1]) - } - return { status: 200 } - } catch (error) { - await SequelizeRepository.rollbackTransaction(transaction) - this.log.error(error, 'Error while adding members to merge') - throw error - } - } - /** * Given two members, add them to the noMerge fields of each other. * @param memberOneId ID of the first member @@ -768,8 +741,9 @@ export default class MemberService extends LoggerBase { try { await MemberRepository.addNoMerge(memberOneId, memberTwoId, txOptions) await MemberRepository.addNoMerge(memberTwoId, memberOneId, txOptions) + + // Removes from either order of the pair await MemberRepository.removeToMerge(memberOneId, memberTwoId, txOptions) - await MemberRepository.removeToMerge(memberTwoId, memberOneId, txOptions) await SequelizeRepository.commitTransaction(transaction) } catch (error) { diff --git a/services/apps/merge_suggestions_worker/src/activities/memberMergeSuggestions.ts b/services/apps/merge_suggestions_worker/src/activities/memberMergeSuggestions.ts index a020e3e8c6..f63ff595fa 100644 --- a/services/apps/merge_suggestions_worker/src/activities/memberMergeSuggestions.ts +++ b/services/apps/merge_suggestions_worker/src/activities/memberMergeSuggestions.ts @@ -2,7 +2,7 @@ import uniqBy from 'lodash.uniqby' import { parseGitHubNoreplyEmail, parseGitLabNoreplyEmail } from '@crowd/common' -import { addMemberNoMerge } from '@crowd/data-access-layer/src/member_merge' +import { insertMemberNoMerge, removeMemberToMerge } from '@crowd/data-access-layer/src/member_merge' import { MemberField, queryMembers } from '@crowd/data-access-layer/src/members' import MemberMergeSuggestionsRepository from '@crowd/data-access-layer/src/old/apps/merge_suggestions_worker/memberMergeSuggestions.repo' import { pgpQx } from '@crowd/data-access-layer/src/queryExecutor' @@ -498,15 +498,14 @@ export async function getRawMemberMergeSuggestions( return memberMergeSuggestionsRepo.getRawMemberSuggestions(similarityFilter, limit) } -export async function removeMemberMergeSuggestion( - suggestion: string[], - table: MemberMergeSuggestionTable, -): Promise { - const memberMergeSuggestionsRepo = new MemberMergeSuggestionsRepository( - svc.postgres.writer.connection(), - svc.log, - ) - await memberMergeSuggestionsRepo.removeMemberMergeSuggestion(suggestion, table) +export async function removeMemberMergeSuggestion(suggestion: string[]): Promise { + if (suggestion.length !== 2) { + svc.log.debug(`Suggestions array must have two ids!`) + return + } + + const qx = pgpQx(svc.postgres.writer.connection()) + await removeMemberToMerge(qx, suggestion[0], suggestion[1]) } export async function addMemberSuggestionToNoMerge(suggestion: string[]): Promise { @@ -516,5 +515,6 @@ export async function addMemberSuggestionToNoMerge(suggestion: string[]): Promis } const qx = pgpQx(svc.postgres.writer.connection()) - await addMemberNoMerge(qx, suggestion[0], suggestion[1]) + await insertMemberNoMerge(qx, suggestion[0], suggestion[1]) + await insertMemberNoMerge(qx, suggestion[1], suggestion[0]) } diff --git a/services/apps/merge_suggestions_worker/src/workflows/mergeMembersWithLLM.ts b/services/apps/merge_suggestions_worker/src/workflows/mergeMembersWithLLM.ts index 236c7f9ea9..40bc18e3f6 100644 --- a/services/apps/merge_suggestions_worker/src/workflows/mergeMembersWithLLM.ts +++ b/services/apps/merge_suggestions_worker/src/workflows/mergeMembersWithLLM.ts @@ -1,6 +1,6 @@ import { continueAsNew, proxyActivities } from '@temporalio/workflow' -import { LLMSuggestionVerdictType, MemberMergeSuggestionTable } from '@crowd/types' +import { LLMSuggestionVerdictType } from '@crowd/types' import type * as activities from '../activities' import { ILLMResult, IProcessMergeMemberSuggestionsWithLLM } from '../types' @@ -78,11 +78,7 @@ export async function mergeMembersWithLLM( if (members.length !== 2) { console.log(`Failed getting members data in suggestion. Skipping suggestion: ${suggestion}`) - await removeMemberMergeSuggestion(suggestion, MemberMergeSuggestionTable.MEMBER_TO_MERGE_RAW) - await removeMemberMergeSuggestion( - suggestion, - MemberMergeSuggestionTable.MEMBER_TO_MERGE_FILTERED, - ) + await removeMemberMergeSuggestion(suggestion) continue } @@ -115,11 +111,7 @@ export async function mergeMembersWithLLM( console.log( `LLM doesn't think these members are the same. Removing from suggestions and adding to no merge: ${suggestion[0]} and ${suggestion[1]}!`, ) - await removeMemberMergeSuggestion( - suggestion, - MemberMergeSuggestionTable.MEMBER_TO_MERGE_FILTERED, - ) - await removeMemberMergeSuggestion(suggestion, MemberMergeSuggestionTable.MEMBER_TO_MERGE_RAW) + await removeMemberMergeSuggestion(suggestion) await addMemberSuggestionToNoMerge(suggestion) } } diff --git a/services/libs/common_services/src/services/member/unmerge.ts b/services/libs/common_services/src/services/member/unmerge.ts index 73c3ae9878..a6055aea8d 100644 --- a/services/libs/common_services/src/services/member/unmerge.ts +++ b/services/libs/common_services/src/services/member/unmerge.ts @@ -21,7 +21,7 @@ import { removeMemberRole, updateMember, } from '@crowd/data-access-layer' -import { addMemberNoMerge } from '@crowd/data-access-layer/src/member_merge' +import { insertMemberNoMerge } from '@crowd/data-access-layer/src/member_merge' import { deleteMemberSegmentAffiliations, findMemberAffiliations, @@ -625,7 +625,7 @@ export async function unmergeMember( } // Add primary and secondary to no merge so they don't get suggested again - await addMemberNoMerge(tx, memberId, secondaryId) + await insertMemberNoMerge(tx, memberId, secondaryId) await setMergeAction(tx, MergeActionType.MEMBER, memberId, secondaryId, { step: MergeActionStep.UNMERGE_SYNC_DONE, 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 7732d7a57c..7f4f9656d5 100644 --- a/services/libs/data-access-layer/src/member_merge/index.ts +++ b/services/libs/data-access-layer/src/member_merge/index.ts @@ -5,37 +5,46 @@ export async function removeMemberToMerge( memberId: string, toMergeId: string, ): Promise { - await qx.result( - ` - DELETE FROM "memberToMerge" - WHERE "memberId" = $(memberId) - AND "toMergeId" = $(toMergeId) - `, - { - memberId, - toMergeId, - }, - ) + const replacements = { memberId, toMergeId } + + const whereClause = ` + WHERE + ("memberId" = $(memberId) AND "toMergeId" = $(toMergeId)) + OR + ("memberId" = $(toMergeId) AND "toMergeId" = $(memberId)) + ` + + await qx.tx(async (tx) => { + await tx.result( + ` + DELETE FROM "memberToMerge" + ${whereClause} + `, + replacements, + ) + + await tx.result( + ` + DELETE FROM "memberToMergeRaw" + ${whereClause} + `, + replacements, + ) + }) } -export async function addMemberNoMerge( +export async function insertMemberNoMerge( qx: QueryExecutor, memberId: string, noMergeId: string, ): Promise { - const currentTime = new Date() await qx.result( ` INSERT INTO "memberNoMerge" ("memberId", "noMergeId", "createdAt", "updatedAt") - VALUES ($(memberId), $(noMergeId), $(createdAt), $(updatedAt)) - on conflict ("memberId", "noMergeId") do nothing + VALUES ($(memberId), $(noMergeId), NOW(), NOW()) + ON CONFLICT ("memberId", "noMergeId") DO NOTHING `, - { - memberId, - noMergeId, - createdAt: currentTime, - updatedAt: currentTime, - }, + { memberId, noMergeId }, ) } diff --git a/services/libs/data-access-layer/src/old/apps/members_enrichment_worker/index.ts b/services/libs/data-access-layer/src/old/apps/members_enrichment_worker/index.ts index e1b8c2bc77..5f0ebed0e2 100644 --- a/services/libs/data-access-layer/src/old/apps/members_enrichment_worker/index.ts +++ b/services/libs/data-access-layer/src/old/apps/members_enrichment_worker/index.ts @@ -365,14 +365,6 @@ export async function findExistingMember( return results.map((r) => r.memberId) } -export async function addMemberToMerge(tx: DbTransaction, memberId: string, toMergeId: string) { - await tx.query( - `INSERT INTO "memberToMerge" ("memberId", "toMergeId", similarity) - VALUES ($1, $2, $3);"`, - [memberId, toMergeId, 0.9], - ) -} - export async function findOrganizationIdentities( tx: DbTransaction, organizationId: string, diff --git a/services/libs/data-access-layer/src/old/apps/merge_suggestions_worker/memberMergeSuggestions.repo.ts b/services/libs/data-access-layer/src/old/apps/merge_suggestions_worker/memberMergeSuggestions.repo.ts index d595c47ee1..0bc5256114 100644 --- a/services/libs/data-access-layer/src/old/apps/merge_suggestions_worker/memberMergeSuggestions.repo.ts +++ b/services/libs/data-access-layer/src/old/apps/merge_suggestions_worker/memberMergeSuggestions.repo.ts @@ -325,31 +325,6 @@ class MemberMergeSuggestionsRepository { return results.map((r) => [r.memberId, r.toMergeId]) } - - async removeMemberMergeSuggestion( - suggestion: string[], - table: MemberMergeSuggestionTable, - ): Promise { - const query = ` - delete from "${table}" - where - ("memberId" = $(memberId) and "toMergeId" = $(toMergeId)) - or - ("memberId" = $(toMergeId) and "toMergeId" = $(memberId)) - ` - - const replacements = { - memberId: suggestion[0], - toMergeId: suggestion[1], - } - - try { - await this.connection.none(query, replacements) - } catch (error) { - this.log.error(`Error removing member suggestions from ${table}`, error) - throw error - } - } } export default MemberMergeSuggestionsRepository