Skip to content

Commit 48c6de2

Browse files
committed
fix: make activity result inserts idempotent on retry (CM-1318)
Signed-off-by: Uroš Marolt <uros@marolt.me>
1 parent 0e43d27 commit 48c6de2

2 files changed

Lines changed: 13 additions & 1 deletion

File tree

services/apps/mailing_list_integration/src/crowdmail/database/crud.py

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -219,9 +219,13 @@ async def release_list(list_id: str) -> None:
219219

220220

221221
async def batch_insert_activities(records: list[tuple], batch_size=100) -> None:
222+
# id is deterministic per (integrationId, shard, git commit) — ON CONFLICT DO
223+
# NOTHING makes a retried batch (after a crash or a Kafka-send failure) idempotent
224+
# instead of inserting duplicate rows for the same messages.
222225
sql_query = """
223226
INSERT INTO integration.results(id, state, data, "tenantId", "integrationId")
224227
values($1, $2, $3, $4, $5)
228+
ON CONFLICT (id) DO NOTHING
225229
"""
226230
logger.info(f"Saving {len(records)} activities into integration.results")
227231
for i in range(0, len(records), batch_size):

services/apps/mailing_list_integration/src/crowdmail/worker/list_worker.py

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -120,7 +120,15 @@ async def _process_single_list(self, mailing_list: MailingList):
120120
continue
121121
activity_data["segmentId"] = mailing_list.segment_id
122122

123-
result_id = str(uuid.uuid1())
123+
# Deterministic (not random) so a retried batch after a crash or a
124+
# Kafka-send failure re-inserts the exact same row instead of a
125+
# duplicate, per commit id which is stable across retries.
126+
result_id = str(
127+
uuid.uuid5(
128+
uuid.NAMESPACE_URL,
129+
f"{mailing_list.integration_id}:{shard}:{git_id}",
130+
)
131+
)
124132
data_dict = {"type": IntegrationResultType.ACTIVITY, "data": activity_data}
125133
activities_db.append(
126134
(

0 commit comments

Comments
 (0)