Skip to content

Commit 83ee8b3

Browse files
committed
Merge branch 'main' into chore/claude-code-setup
2 parents 503f112 + 20f8fe4 commit 83ee8b3

17 files changed

Lines changed: 345 additions & 46 deletions

File tree

backend/src/database/flyway_migrate.sh

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ flyway \
1010
-password="$PGPASSWORD" \
1111
-connectRetries=60 \
1212
-outOfOrder=true \
13+
-mixed=true \
1314
-placeholderReplacement=false \
1415
-schemas=public \
1516
-X \

backend/src/database/migrations/U1774609007__data-sink-worker-optimizations.sql

Whitespace-only changes.
Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,11 @@
1+
-- Drop 4 unused activityRelations indexes (already dropped on prod 2026-03-27,
2+
-- see ACTIVITYRELATIONS_INDEX_CLEANUP.md — IF EXISTS guards for idempotency)
3+
alter table "activityRelations" drop constraint if exists "activityRelations_activityId_memberId_key";
4+
5+
drop index concurrently if exists "ix_activityRelations_memberId_segmentId_include";
6+
drop index concurrently if exists "ix_activityRelations_organizationId_segmentId_include";
7+
drop index concurrently if exists "ix_activityRelations_platform_username";
8+
9+
create index concurrently if not exists idx_osa_org_segment_membercount
10+
on "organizationSegmentsAgg" ("organizationId", "segmentId")
11+
include ("memberCount");

backend/src/product/flyway_migrate.sh

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ flyway \
1010
-password="$PGPASSWORD" \
1111
-connectRetries=60 \
1212
-outOfOrder=true \
13+
-mixed=true \
1314
-placeholderReplacement=false \
1415
-schemas=public \
1516
-X \

services/apps/data_sink_worker/src/service/activity.service.ts

Lines changed: 40 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -28,7 +28,6 @@ import {
2828
} from '@crowd/data-access-layer'
2929
import { IDbActivityRelation } from '@crowd/data-access-layer/src/activityRelations/types'
3030
import { DbStore, arePrimitivesDbEqual } from '@crowd/data-access-layer/src/database'
31-
import { getMemberNoMerge } from '@crowd/data-access-layer/src/member_merge'
3231
import {
3332
IActivityRelationCreateOrUpdateData,
3433
IDbActivity,
@@ -56,7 +55,7 @@ import {
5655
} from '@crowd/types'
5756

5857
import { IActivityUpdateData, ISentimentActivityInput } from './activity.data'
59-
import MemberService from './member.service'
58+
import MemberService, { mergeIfAllowed } from './member.service'
6059
import { IProcessActivityResult } from './types'
6160

6261
/* eslint-disable @typescript-eslint/no-explicit-any */
@@ -293,6 +292,28 @@ export default class ActivityService extends LoggerBase {
293292
}
294293
}
295294

295+
// When activity.username is set but differs from the member's platform identity value,
296+
// override it so the member lookup and the identity insert use the same key.
297+
// Example: git activities set activity.username to the author display name (e.g. "John Doe")
298+
// while the identity stores the email (e.g. "john.doe@example.com"). Without this correction
299+
// the lookup misses the existing member, creating an unnecessary orphan member.
300+
if (username && member) {
301+
const platformIdentity = member.identities.find(
302+
(i) =>
303+
i.platform === platform &&
304+
i.type === MemberIdentityType.USERNAME &&
305+
i.value &&
306+
i.verified,
307+
)
308+
if (platformIdentity && platformIdentity.value !== username) {
309+
this.log.debug(
310+
{ platform, originalUsername: username, correctedUsername: platformIdentity.value },
311+
'Overriding activity.username with member platform identity value',
312+
)
313+
activity.username = platformIdentity.value
314+
}
315+
}
316+
296317
member.identities = member.identities.filter((i) => i.value)
297318

298319
if (!username) {
@@ -1721,32 +1742,24 @@ export default class ActivityService extends LoggerBase {
17211742
const originalId = metadata.memberWithIdentity as string
17221743
const targetId = metadata.memberIdToUpdate as string
17231744

1724-
// but first check memberNoMerge table
1725-
const noMergeMemberIds = await getMemberNoMerge(this.pgQx, [originalId, targetId])
1726-
1727-
const noMerge = singleOrDefault(
1728-
noMergeMemberIds,
1729-
(m) =>
1730-
(m.memberId === originalId && m.noMergeId === targetId) ||
1731-
(m.memberId === targetId && m.noMergeId === originalId),
1732-
)
1733-
1734-
if (noMerge) {
1735-
metadata.noMerge = true
1736-
} else {
1737-
try {
1738-
await this.pgQx.tx(async (txPgQx) => {
1739-
const service = new CommonMemberService(txPgQx, this.temporal, this.log)
1740-
await service.merge(originalId, targetId)
1741-
})
1742-
1745+
try {
1746+
const merged = await mergeIfAllowed(
1747+
this.pgQx,
1748+
this.temporal,
1749+
this.log,
1750+
originalId,
1751+
targetId,
1752+
)
1753+
if (merged) {
17431754
return originalId
1744-
} catch (err) {
1745-
metadata.mergeError = {
1746-
errorMessage: err?.message ?? '<no error message>',
1747-
errorStack: err?.stack,
1748-
err,
1749-
}
1755+
} else {
1756+
metadata.noMerge = true
1757+
}
1758+
} catch (err) {
1759+
metadata.mergeError = {
1760+
errorMessage: err?.message ?? '<no error message>',
1761+
errorStack: err?.stack,
1762+
err,
17501763
}
17511764
}
17521765
}

0 commit comments

Comments
 (0)