Skip to content

Commit 9ff36e7

Browse files
committed
fix: some refactoring and better control on logExecutionTimeV2 duration logging
Signed-off-by: Uroš Marolt <uros@marolt.me>
1 parent a8c8823 commit 9ff36e7

4 files changed

Lines changed: 39 additions & 26 deletions

File tree

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

Lines changed: 22 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -1581,20 +1581,6 @@ export default class ActivityService extends LoggerBase {
15811581
)
15821582
}
15831583

1584-
await this.searchSyncWorkerEmitter.triggerMemberSync(
1585-
prepared.payload.memberId,
1586-
onboarding,
1587-
prepared.payload.segmentId,
1588-
)
1589-
1590-
if (prepared.payload.objectMemberId) {
1591-
await this.searchSyncWorkerEmitter.triggerMemberSync(
1592-
prepared.payload.objectMemberId,
1593-
onboarding,
1594-
prepared.payload.segmentId,
1595-
)
1596-
}
1597-
15981584
if (prepared.payload.organizationId) {
15991585
await this.redisClient.sAdd(
16001586
'organizationIdsForAggComputation',
@@ -1603,6 +1589,28 @@ export default class ActivityService extends LoggerBase {
16031589
}
16041590
}
16051591

1592+
// Deduplicate member sync triggers — a member may appear in many activities in the
1593+
// same batch. Emit once per unique (memberId, segmentId) pair.
1594+
const memberSyncKeys = new Set<string>()
1595+
const memberSyncPromises: Promise<void>[] = []
1596+
for (const prepared of preparedForUpsert) {
1597+
for (const memberId of [prepared.payload.memberId, prepared.payload.objectMemberId]) {
1598+
if (!memberId) continue
1599+
const key = `${memberId}:${prepared.payload.segmentId}`
1600+
if (!memberSyncKeys.has(key)) {
1601+
memberSyncKeys.add(key)
1602+
memberSyncPromises.push(
1603+
this.searchSyncWorkerEmitter.triggerMemberSync(
1604+
memberId,
1605+
onboarding,
1606+
prepared.payload.segmentId,
1607+
),
1608+
)
1609+
}
1610+
}
1611+
}
1612+
await Promise.all(memberSyncPromises)
1613+
16061614
for (const prepared of preparedActivities) {
16071615
resultMap.set(prepared.resultId, { success: true })
16081616
}

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

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -229,7 +229,6 @@ export default class MemberService extends LoggerBase {
229229

230230
if (organizations.length > 0) {
231231
const uniqOrgs = uniqby(organizations, 'id')
232-
const orgService = new OrganizationService(this.store, this.log)
233232

234233
const orgsToAdd = (
235234
await Promise.all(
@@ -451,7 +450,6 @@ export default class MemberService extends LoggerBase {
451450

452451
if (organizations.length > 0) {
453452
const uniqOrgs = uniqby(organizations, 'id')
454-
const orgService = new OrganizationService(this.store, this.log)
455453

456454
this.log.trace({ memberId: id }, 'Finding member organizations!')
457455
const orgsToAdd = (

services/libs/database/src/connection.ts

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -129,7 +129,9 @@ export const getDbConnection = async (
129129
})
130130

131131
const profile = process.env['CROWD_POSTGRESQL_PROFILE_QUERIES'] !== undefined
132-
const minQueryDuration = Number(process.env['CROWD_POSTGRESQL_PROFILE_QUERIES_MIN_DURATION'] || 0)
132+
const minQueryDurationMs = Number(
133+
process.env['CROWD_POSTGRESQL_PROFILE_QUERIES_MIN_DURATION'] || 0,
134+
)
133135

134136
const oldQuery = client.query
135137
// eslint-disable-next-line @typescript-eslint/no-explicit-any
@@ -141,7 +143,7 @@ export const getDbConnection = async (
141143
return result
142144
} finally {
143145
const duration = performance.now() - start
144-
if (profile && duration >= minQueryDuration) {
146+
if (profile && duration >= minQueryDurationMs) {
145147
const durationSeconds = duration / 1000.0
146148
log.warn(
147149
{ durationSeconds: durationSeconds.toFixed(2), query, values: options },

services/libs/logging/src/utility.ts

Lines changed: 13 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -27,24 +27,29 @@ export const logExecutionTimeV2 = async <T>(
2727
if (!process.env.CROWD_LOG_EXECUTION_TIME) {
2828
return toProcess()
2929
}
30+
31+
const minDurationMs = Number(process.env['CROWD_LOG_EXECUTION_TIME_MIN_DURATION'] || 0)
32+
3033
const start = performance.now()
3134

32-
const end = () => {
33-
const end = performance.now()
34-
const duration = end - start
35-
const durationInSeconds = duration / 1000
36-
return durationInSeconds.toFixed(2)
37-
}
35+
const formatDuration = (durationMs: number) => (durationMs / 1000).toFixed(2)
36+
3837
try {
3938
if (process.env.CROWD_LOG_EXECUTION_START) {
4039
log.info(`Starting process ${name}...`)
4140
}
4241

4342
const result = await toProcess()
44-
log.info(`Process ${name} took ${end()} seconds!`)
43+
const durationMs = performance.now() - start
44+
if (durationMs >= minDurationMs) {
45+
log.info(`Process ${name} took ${formatDuration(durationMs)} seconds!`)
46+
}
4547
return result
4648
} catch (e) {
47-
log.info(`Process ${name} failed after ${end()} seconds!`)
49+
const durationMs = performance.now() - start
50+
if (durationMs >= minDurationMs) {
51+
log.info(`Process ${name} failed after ${formatDuration(durationMs)} seconds!`)
52+
}
4853
throw e
4954
}
5055
}

0 commit comments

Comments
 (0)