Skip to content

Commit f29db8d

Browse files
authored
fix: split dissectMember backup fetch to avoid temporal payload limits (#4288)
Signed-off-by: Yeganathan S <63534555+skwowet@users.noreply.github.com>
1 parent 673509b commit f29db8d

4 files changed

Lines changed: 47 additions & 2 deletions

File tree

services/apps/script_executor_worker/src/activities.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,7 @@ import {
3737
findMemberById,
3838
findMemberIdentitiesGroupedByPlatform,
3939
findMemberMergeActions,
40+
findMergeActionUnmergeBackup,
4041
} from './activities/dissect-member'
4142
import {
4243
getBotMembersWithOrgAffiliation,
@@ -66,6 +67,7 @@ export {
6667
findMembersWithSamePlatformIdentitiesDifferentCapitalization,
6768
mergeMembers,
6869
findMemberMergeActions,
70+
findMergeActionUnmergeBackup,
6971
unmergeMembers,
7072
unmergeMembersPreview,
7173
waitForTemporalWorkflowExecutionFinish,

services/apps/script_executor_worker/src/activities/dissect-member/index.ts

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,17 @@ export async function findMemberMergeActions(
3030
return mergeActions
3131
}
3232

33+
export async function findMergeActionUnmergeBackup(
34+
mergeActionId: string,
35+
): Promise<IMergeAction['unmergeBackup'] | null> {
36+
try {
37+
const mergeActionRepo = new MergeActionRepository(svc.postgres.reader.connection(), svc.log)
38+
return await mergeActionRepo.findMergeActionUnmergeBackup(mergeActionId)
39+
} catch (err) {
40+
throw new Error(err)
41+
}
42+
}
43+
3344
export async function findMemberIdentitiesGroupedByPlatform(
3445
memberId: string,
3546
): Promise<IFindMemberIdentitiesGroupedByPlatformResult[]> {

services/apps/script_executor_worker/src/workflows/dissectMember.ts

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -89,9 +89,18 @@ export async function dissectMember(args: IDissectMemberArgs): Promise<void> {
8989
// 2. wait for temporal async stuff to complete
9090
// 3. call the same workflow again for the unmerged secondary member
9191

92+
// Fetch each backup in its own activity — listing 10 backups in one result can exceed
93+
// Temporal's 2MB activity payload limit on polluted members.
94+
const unmergeBackup = await activity.findMergeActionUnmergeBackup(mergeAction.id)
95+
96+
if (!unmergeBackup) {
97+
console.log(`Merge action ${mergeAction.id} has no unmerge backup, skipping!`)
98+
continue
99+
}
100+
92101
await common.unmergeMembers(
93102
mergeAction.primaryId,
94-
mergeAction.unmergeBackup as IUnmergeBackup<IMemberUnmergeBackup>,
103+
unmergeBackup as IUnmergeBackup<IMemberUnmergeBackup>,
95104
)
96105

97106
const workflowId = `finishMemberUnmerging/${mergeAction.primaryId}/${mergeAction.secondaryId}`

services/libs/data-access-layer/src/old/apps/script_executor_worker/mergeAction.repo.ts

Lines changed: 24 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,16 @@ class MergeActionRepository {
1919
) {
2020
let rows: IMergeAction[] = []
2121
let query = `
22-
select * from "mergeActions" ma
22+
select
23+
ma.id,
24+
ma.type,
25+
ma."primaryId",
26+
ma."secondaryId",
27+
ma."createdAt",
28+
ma."updatedAt",
29+
ma.step,
30+
ma.state
31+
from "mergeActions" ma
2332
where
2433
ma."state" = 'merged' and
2534
ma."primaryId" = $(memberId) and
@@ -62,6 +71,20 @@ class MergeActionRepository {
6271
return rows
6372
}
6473

74+
async findMergeActionUnmergeBackup(mergeActionId: string) {
75+
const rows = await this.connection.query(
76+
`
77+
select ma."unmergeBackup"
78+
from "mergeActions" ma
79+
where ma.id = $(mergeActionId)
80+
and ma."unmergeBackup" is not null
81+
`,
82+
{ mergeActionId },
83+
)
84+
85+
return rows[0]?.unmergeBackup ?? null
86+
}
87+
6588
async findMergeActionsWithDeletedSecondaryEntities(
6689
limit: number,
6790
offset: number,

0 commit comments

Comments
 (0)