Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion services/apps/entity_merging_worker/src/activities.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ export {
notifyFrontendOrganizationMergeSuccessful,
notifyFrontendOrganizationUnmergeSuccessful,
syncOrganization,
recalculateActivityAffiliationsOfOrganizationSynchronous,
recalculateActivityAffiliationsOfOrganizationAsync,
} from './activities/organizations'

export { setMergeAction } from './activities/common'
6 changes: 3 additions & 3 deletions services/apps/entity_merging_worker/src/activities/members.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { WorkflowIdReusePolicy } from '@temporalio/workflow'
import { WorkflowIdConflictPolicy, WorkflowIdReusePolicy } from '@temporalio/workflow'

import { DEFAULT_TENANT_ID } from '@crowd/common'
import {
Expand Down Expand Up @@ -47,8 +47,8 @@ export async function recalculateActivityAffiliationsOfMemberAsync(
await svc.temporal.workflow.start('memberUpdate', {
taskQueue: 'profiles',
workflowId,
// if the workflow is already running, this policy will throw an error
workflowIdReusePolicy: WorkflowIdReusePolicy.REJECT_DUPLICATE,
workflowIdReusePolicy: WorkflowIdReusePolicy.ALLOW_DUPLICATE,
workflowIdConflictPolicy: WorkflowIdConflictPolicy.FAIL,
retry: {
maximumAttempts: 10,
},
Expand Down
68 changes: 47 additions & 21 deletions services/apps/entity_merging_worker/src/activities/organizations.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
import { WorkflowIdConflictPolicy, WorkflowIdReusePolicy } from '@temporalio/workflow'

import { DEFAULT_TENANT_ID } from '@crowd/common'
import { moveActivityRelationsToAnotherOrganization } from '@crowd/data-access-layer/src/activityRelations'
import {
Expand Down Expand Up @@ -35,31 +37,55 @@ export async function finishOrganizationMergingUpdateActivities(
await moveActivityRelationsToAnotherOrganization(qx, secondaryId, primaryId)
}

export async function recalculateActivityAffiliationsOfOrganizationSynchronous(
export async function recalculateActivityAffiliationsOfOrganizationAsync(
organizationId: string,
): Promise<void> {
await svc.temporal.workflow.start('organizationUpdate', {
taskQueue: 'profiles',
workflowId: `${TemporalWorkflowId.ORGANIZATION_UPDATE}/${organizationId}`,
followRuns: true,
retry: {
maximumAttempts: 10,
},
args: [
{
organization: {
id: organizationId,
},
recalculateAffiliations: true,
syncOptions: {
doSync: false,
withAggs: false,
},
const workflowId = `${TemporalWorkflowId.ORGANIZATION_UPDATE}/${organizationId}`

try {
const handle = svc.temporal.workflow.getHandle(workflowId)
const { status } = await handle.describe()

if (status.name === 'RUNNING') {
await handle.result()
}
} catch (err) {
if (err.name !== 'WorkflowNotFoundError') {
svc.log.error({ err }, 'Failed to check workflow state')
throw err
}
}

try {
await svc.temporal.workflow.start('organizationUpdate', {
taskQueue: 'profiles',
workflowId,
workflowIdReusePolicy: WorkflowIdReusePolicy.ALLOW_DUPLICATE,
workflowIdConflictPolicy: WorkflowIdConflictPolicy.FAIL,
retry: {
maximumAttempts: 10,
},
],
})
args: [
{
organization: {
id: organizationId,
},
recalculateAffiliations: true,
syncOptions: {
doSync: false,
withAggs: false,
},
},
],
})
} catch (err) {
if (err.name === 'WorkflowExecutionAlreadyStartedError') {
svc.log.info({ workflowId }, 'Workflow already started, skipping')
return
}

await svc.temporal.workflow.result(`${TemporalWorkflowId.ORGANIZATION_UPDATE}/${organizationId}`)
throw err
}
}

export async function syncOrganization(organizationId: string, syncStart: Date): Promise<void> {
Expand Down
6 changes: 3 additions & 3 deletions services/apps/entity_merging_worker/src/workflows/all.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ const {
notifyFrontendOrganizationMergeSuccessful,
notifyFrontendOrganizationUnmergeSuccessful,
recalculateActivityAffiliationsOfMemberAsync,
recalculateActivityAffiliationsOfOrganizationSynchronous,
recalculateActivityAffiliationsOfOrganizationAsync,
setMergeAction,
syncMember,
syncOrganization,
Expand Down Expand Up @@ -118,8 +118,8 @@ export async function finishOrganizationUnmerging(
await setMergeAction(primaryId, secondaryId, {
step: MergeActionStep.UNMERGE_ASYNC_STARTED,
})
await recalculateActivityAffiliationsOfOrganizationSynchronous(primaryId)
await recalculateActivityAffiliationsOfOrganizationSynchronous(secondaryId)
await recalculateActivityAffiliationsOfOrganizationAsync(primaryId)
await recalculateActivityAffiliationsOfOrganizationAsync(secondaryId)
const syncStart = new Date()
await syncOrganization(primaryId, syncStart)
await syncOrganization(secondaryId, syncStart)
Expand Down