diff --git a/services/apps/cron_service/src/jobs/inferMemberOrganizationStintChanges.job.ts b/services/apps/cron_service/src/jobs/inferMemberOrganizationStintChanges.job.ts index 369fc4d6bb..9cc930b7ff 100644 --- a/services/apps/cron_service/src/jobs/inferMemberOrganizationStintChanges.job.ts +++ b/services/apps/cron_service/src/jobs/inferMemberOrganizationStintChanges.job.ts @@ -15,14 +15,21 @@ import { fetchManyOrganizationAffiliationPolicies, fetchMemberOrganizationsBySource, findMemberById, + findOrgsByIds, updateMemberOrganization, } from '@crowd/data-access-layer' import { WRITE_DB_CONFIG, getDbConnection } from '@crowd/data-access-layer/src/database' import { deleteMemberSegmentAffiliations } from '@crowd/data-access-layer/src/member_segment_affiliations' import { pgpQx } from '@crowd/data-access-layer/src/queryExecutor' +import { Logger } from '@crowd/logging' import { REDIS_CONFIG, RedisCache, RedisClient, getRedisClient } from '@crowd/redis' import { TEMPORAL_CONFIG, getTemporalClient } from '@crowd/temporal' -import { MemberOrgDate, MemberOrgStintChange, OrganizationSource } from '@crowd/types' +import { + IMemberOrganization, + MemberOrgDate, + MemberOrgStintChange, + OrganizationSource, +} from '@crowd/types' import { IJobDefinition } from '../types' @@ -78,7 +85,15 @@ const job: IJobDefinition = { { withDeleted: true }, ) - const changes = inferMemberOrganizationStintChanges(memberId, existingOrgs, orgDates) + const validOrgDates = await dropStaleOrganizationDates( + qx, + existingOrgs, + orgDates, + memberId, + ctx.log, + ) + + const changes = inferMemberOrganizationStintChanges(memberId, existingOrgs, validOrgDates) if (changes.length > 0) { ctx.log.debug({ memberId, changes }, 'Stint changes identified.') @@ -104,6 +119,11 @@ const job: IJobDefinition = { processed++ } catch (err) { + if ((err as { code?: string })?.code === '23503') { + ctx.log.warn(err, { memberId }, 'Stint change referenced missing organization.') + continue + } + ctx.log.error(err, { memberId }, 'Failed to process member stint inference.') throw err } @@ -129,6 +149,42 @@ function parseSetMembers(members: string[]): MemberOrgDate[] { return results } +/** + * Drops queued dates for organizations that no longer exist (deleted/merged after being + * queued), so a single stale reference can't FK-violate and roll back the whole member's + * transaction, taking other, still-valid queued dates down with it. + */ +async function dropStaleOrganizationDates( + qx: QueryExecutor, + existingOrgs: IMemberOrganization[], + orgDates: MemberOrgDate[], + memberId: string, + log: Logger, +): Promise { + const knownOrgIds = new Set(existingOrgs.map((o) => o.organizationId)) + const orgIdsToVerify = [...new Set(orgDates.map((d) => d.organizationId))].filter( + (id) => !knownOrgIds.has(id), + ) + + if (orgIdsToVerify.length === 0) { + return orgDates + } + + const existingOrgIds = new Set((await findOrgsByIds(qx, orgIdsToVerify)).map((o) => o.id)) + const staleOrgIds = orgIdsToVerify.filter((id) => !existingOrgIds.has(id)) + + if (staleOrgIds.length === 0) { + return orgDates + } + + log.warn( + { memberId, staleOrgIds }, + 'Dropping queued stint dates for organizations that no longer exist.', + ) + + return orgDates.filter((d) => !staleOrgIds.includes(d.organizationId)) +} + /** * Purges a member from the queue and their associated Redis entries. */ diff --git a/services/libs/common_services/src/services/member-organization.ts b/services/libs/common_services/src/services/member-organization.ts index 5c2d9083f1..7e53b4302b 100644 --- a/services/libs/common_services/src/services/member-organization.ts +++ b/services/libs/common_services/src/services/member-organization.ts @@ -175,6 +175,10 @@ export function inferMemberOrganizationStintChanges( const activeRows = normalizedRows.filter((row) => !row.deletedAt) + const tombstonedOrgIds = new Set( + normalizedRows.filter((row) => row.deletedAt && row.deletedBy).map((row) => row.organizationId), + ) + // Deleted dated rows suppress recreation for dates the user removed const deletedRows = normalizedRows.filter( (row): row is typeof row & { dateStart: string } => !!row.deletedAt && !!row.dateStart, @@ -199,6 +203,10 @@ export function inferMemberOrganizationStintChanges( })) for (const { organizationId, date: targetDate } of sortedDates) { + if (tombstonedOrgIds.has(organizationId)) { + continue + } + if ( deletedRows.some( (row) => diff --git a/services/libs/data-access-layer/src/members/organizations.ts b/services/libs/data-access-layer/src/members/organizations.ts index e8b1f46432..de6346fe57 100644 --- a/services/libs/data-access-layer/src/members/organizations.ts +++ b/services/libs/data-access-layer/src/members/organizations.ts @@ -83,7 +83,8 @@ export async function fetchMemberOrganizationsBySource( "title", "memberId", "source", - "deletedAt" + "deletedAt", + "deletedBy" FROM "memberOrganizations" WHERE "memberId" = $(memberId) AND "source" = $(source) diff --git a/services/libs/types/src/organizations.ts b/services/libs/types/src/organizations.ts index 78e74f4ed5..abd2dd5ead 100644 --- a/services/libs/types/src/organizations.ts +++ b/services/libs/types/src/organizations.ts @@ -63,6 +63,7 @@ export interface IMemberOrganization { verified?: boolean verifiedBy?: string deletedAt?: string + deletedBy?: string displayName?: string affiliationOverride?: IMemberOrganizationAffiliationOverride }