Skip to content

Commit ec60e44

Browse files
committed
fix: comments
Signed-off-by: Uroš Marolt <uros@marolt.me>
1 parent 84493c7 commit ec60e44

3 files changed

Lines changed: 19 additions & 13 deletions

File tree

services/apps/cron_service/src/jobs/incomingWebhooksCheck.job.ts

Lines changed: 2 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -70,17 +70,9 @@ const job: IJobDefinition = {
7070
const emitter = new IntegrationStreamWorkerEmitter(queueService, ctx.log)
7171
await emitter.init()
7272

73-
let triggered = 0
74-
for (const webhook of webhooks) {
75-
await emitter.triggerWebhookProcessing(webhook.platform, webhook.id)
76-
triggered++
77-
78-
if (triggered % 100 === 0) {
79-
ctx.log.info(`Re-triggered ${triggered} webhooks!`)
80-
}
81-
}
73+
await emitter.triggerWebhookProcessingBatch(webhooks.map((w) => w.id))
8274

83-
ctx.log.info(`Re-triggered ${triggered} stuck pending webhooks in total!`)
75+
ctx.log.info(`Re-triggered ${webhooks.length} stuck pending webhooks in total!`)
8476
},
8577
}
8678

services/libs/common_services/src/services/emitters/integrationStreamWorker.emitter.ts

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import { generateUUIDv1 } from '@crowd/common'
1+
import { generateUUIDv1, partition } from '@crowd/common'
22
import { Logger } from '@crowd/logging'
33
import { CrowdQueue, IQueue } from '@crowd/queue'
44
import {
@@ -59,4 +59,16 @@ export class IntegrationStreamWorkerEmitter extends QueuePriorityService {
5959
{ onboarding: true },
6060
)
6161
}
62+
63+
public async triggerWebhookProcessingBatch(webhookIds: string[]): Promise<void> {
64+
for (const batch of partition(webhookIds, 10)) {
65+
await this.sendMessages(
66+
batch.map((webhookId) => ({
67+
payload: new ProcessWebhookStreamQueueMessage(webhookId),
68+
groupId: generateUUIDv1(),
69+
deduplicationId: webhookId,
70+
})),
71+
)
72+
}
73+
}
6274
}

services/libs/data-access-layer/src/old/apps/data_sink_worker/repo/requestedForErasureMemberIdentities.repo.ts

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -74,10 +74,11 @@ export default class RequestedForErasureMemberIdentitiesRepository extends Repos
7474
}
7575

7676
if (identity.type === MemberIdentityType.EMAIL) {
77+
// SQL matches case-insensitively (lower(value) = lower(?)), so mirror that here.
7778
const row = singleOrDefault(data, (r) => {
7879
return (
7980
r.type === identity.type &&
80-
r.value === identity.value &&
81+
r.value.toLowerCase() === identity.value.toLowerCase() &&
8182
r.platform === identity.platform
8283
)
8384
})
@@ -93,11 +94,12 @@ export default class RequestedForErasureMemberIdentitiesRepository extends Repos
9394
// (type, value) but on different platforms would both match. The previous
9495
// singleOrDefault call threw "Array contains more than one matching element!" in that
9596
// case — a deterministic crash that never self-heals.
97+
// SQL matches case-insensitively (lower(value) = lower(?)), so mirror that here.
9698
const row =
9799
data.find(
98100
(r) =>
99101
r.type === identity.type &&
100-
r.value === identity.value &&
102+
r.value.toLowerCase() === identity.value.toLowerCase() &&
101103
r.platform === identity.platform,
102104
) ?? null
103105

0 commit comments

Comments
 (0)