diff --git a/apps/api/src/services/channel-account-purge.integration.test.ts b/apps/api/src/services/channel-account-purge.integration.test.ts index ba87cdc6..0bd66418 100644 --- a/apps/api/src/services/channel-account-purge.integration.test.ts +++ b/apps/api/src/services/channel-account-purge.integration.test.ts @@ -17,61 +17,104 @@ import { const integrationTest = process.env.RUN_DB_INTEGRATION === "1" ? test : test.skip; +async function provisionTenant(label: string) { + const companyId = crypto.randomUUID(); + const schemaName = getSchemaName(companyId); + const ownerId = crypto.randomUUID(); + await db + .insertInto("users") + .values({ + id: ownerId, + email: `${label}-${ownerId}@example.com`, + password_hash: "test", + }) + .execute(); + await db + .insertInto("companies") + .values({ + id: companyId, + name: `${label} test`, + schema_name: schemaName, + status: "active", + }) + .execute(); + await db + .insertInto("sla_policies") + .values({ + company_id: companyId, + target_minutes: 60, + direct_resolution_target_minutes: 480, + group_response_target_minutes: 120, + group_resolution_target_minutes: 960, + timezone: "UTC", + weekly_schedule: JSON.stringify(DEFAULT_SLA_WEEKLY_SCHEDULE), + exceptions: JSON.stringify([]), + effective_from: new Date("1970-01-01T00:00:00Z"), + created_by: ownerId, + }) + .execute(); + await createTenantSchema(companyId); + await reconcileChannelSpineConcurrentIndexes(db, schemaName); + const tenantDb = getTenantConnection(companyId); + return { companyId, schemaName, ownerId, tenantDb }; +} + +async function teardownTenant( + companyId: string, + schemaName: string, + ownerId: string, +) { + await clearTenantConnection(companyId); + await sql.raw(`DROP SCHEMA IF EXISTS "${schemaName}" CASCADE`).execute(db); + await db + .deleteFrom("sla_policies") + .where("company_id", "=", companyId) + .execute(); + await db.deleteFrom("companies").where("id", "=", companyId).execute(); + await db.deleteFrom("users").where("id", "=", ownerId).execute(); +} + +async function createArchivedTelegramAccount( + tenantDb: ReturnType, + accountId: string, +) { + await tenantDb + .insertInto("channel_accounts") + .values({ + id: accountId, + channel: "telegram", + provider: "telegram_bot", + display_name: "Bot", + status: "archived", + archived_at: new Date(), + }) + .execute(); +} + +async function reactionCountByAccount( + tenantDb: ReturnType, + accountId: string, +) { + return Number( + ( + await tenantDb + .selectFrom("message_reactions") + .select((eb) => eb.fn.countAll().as("count")) + .where("channel_account_id", "=", accountId) + .executeTakeFirstOrThrow() + ).count, + ); +} + describe("purgeArchivedChannelAccount", () => { integrationTest( - "purges a Telegram account whose conversations have no contact", + "purges an archived Telegram account along with its reactions", async () => { - const companyId = crypto.randomUUID(); - const schemaName = getSchemaName(companyId); - const ownerId = crypto.randomUUID(); + const { companyId, schemaName, ownerId, tenantDb } = + await provisionTenant("channel-purge"); try { - await db - .insertInto("users") - .values({ - id: ownerId, - email: `channel-purge-${ownerId}@example.com`, - password_hash: "test", - }) - .execute(); - await db - .insertInto("companies") - .values({ - id: companyId, - name: "Channel purge test", - schema_name: schemaName, - status: "active", - }) - .execute(); - await db - .insertInto("sla_policies") - .values({ - company_id: companyId, - target_minutes: 60, - direct_resolution_target_minutes: 480, - group_response_target_minutes: 120, - group_resolution_target_minutes: 960, - timezone: "UTC", - weekly_schedule: JSON.stringify(DEFAULT_SLA_WEEKLY_SCHEDULE), - exceptions: JSON.stringify([]), - effective_from: new Date("1970-01-01T00:00:00Z"), - created_by: ownerId, - }) - .execute(); - await createTenantSchema(companyId); - await reconcileChannelSpineConcurrentIndexes(db, schemaName); - const tenantDb = getTenantConnection(companyId); const accountId = crypto.randomUUID(); - await tenantDb - .insertInto("channel_accounts") - .values({ - id: accountId, - channel: "telegram", - provider: "telegram_bot", - display_name: "Bot", - status: "archived", - archived_at: new Date(), - }) - .execute(); + await createArchivedTelegramAccount(tenantDb, accountId); const conversation = await tenantDb .insertInto("conversations") .values({ @@ -83,8 +126,9 @@ describe("purgeArchivedChannelAccount", () => { }) .returning("id") .executeTakeFirstOrThrow(); + const messageId = crypto.randomUUID(); + const reactionId = crypto.randomUUID(); await tenantDb.transaction().execute(async (trx) => { - const messageId = crypto.randomUUID(); await trx .insertInto("messages") .values({ @@ -98,6 +142,16 @@ describe("purgeArchivedChannelAccount", () => { timestamp: new Date(), }) .execute(); + await trx + .insertInto("message_reactions") + .values({ + id: reactionId, + message_id: messageId, + reactor_jid: `telegram:${crypto.randomUUID()}`, + emoji: "👍", + channel_account_id: accountId, + }) + .execute(); await openOrReopenCaseForInboundConversation( trx, companyId, @@ -127,6 +181,14 @@ describe("purgeArchivedChannelAccount", () => { .where("id", "=", conversation.id) .executeTakeFirst(), ).toBeUndefined(); + expect( + await tenantDb + .selectFrom("message_reactions") + .select("id") + .where("id", "=", reactionId) + .executeTakeFirst(), + ).toBeUndefined(); + expect(await reactionCountByAccount(tenantDb, accountId)).toBe(0); expect( Number( ( @@ -138,16 +200,75 @@ describe("purgeArchivedChannelAccount", () => { ), ).toBe(0); } finally { - await clearTenantConnection(companyId); - await sql - .raw(`DROP SCHEMA IF EXISTS "${schemaName}" CASCADE`) - .execute(db); - await db - .deleteFrom("sla_policies") - .where("company_id", "=", companyId) + await teardownTenant(companyId, schemaName, ownerId); + } + }, + ); + + integrationTest( + "reaps orphaned reactions whose message row was already deleted", + async () => { + const { companyId, schemaName, ownerId, tenantDb } = + await provisionTenant("channel-purge-orphan"); + try { + const accountId = crypto.randomUUID(); + await createArchivedTelegramAccount(tenantDb, accountId); + const conversation = await tenantDb + .insertInto("conversations") + .values({ + channel_account_id: accountId, + client_thread_key: `telegram:${crypto.randomUUID()}`, + kind: "direct", + subject: "Orphan", + legacy_contact_id: null, + }) + .returning("id") + .executeTakeFirstOrThrow(); + const realMessageId = crypto.randomUUID(); + await tenantDb + .insertInto("messages") + .values({ + id: realMessageId, + contact_id: null, + conversation_id: conversation.id, + channel_account_id: accountId, + from_me: false, + message_type: "text", + content: "hello", + timestamp: new Date(), + }) .execute(); - await db.deleteFrom("companies").where("id", "=", companyId).execute(); - await db.deleteFrom("users").where("id", "=", ownerId).execute(); + const orphanedReactionId = crypto.randomUUID(); + await tenantDb + .insertInto("message_reactions") + .values({ + id: orphanedReactionId, + message_id: crypto.randomUUID(), + reactor_jid: `telegram:${crypto.randomUUID()}`, + emoji: "❤️", + channel_account_id: accountId, + }) + .execute(); + + const result = await purgeArchivedChannelAccount(tenantDb, accountId); + expect(result.deletedMessageCount).toBe(1); + expect( + await tenantDb + .selectFrom("channel_accounts") + .select("id") + .where("id", "=", accountId) + .executeTakeFirst(), + ).toBeUndefined(); + expect( + await tenantDb + .selectFrom("message_reactions") + .select("id") + .where("id", "=", orphanedReactionId) + .executeTakeFirst(), + ).toBeUndefined(); + expect(await reactionCountByAccount(tenantDb, accountId)).toBe(0); + } finally { + await teardownTenant(companyId, schemaName, ownerId); } }, ); diff --git a/apps/api/src/services/channel-account-purge.service.ts b/apps/api/src/services/channel-account-purge.service.ts index 12748a59..0934776e 100644 --- a/apps/api/src/services/channel-account-purge.service.ts +++ b/apps/api/src/services/channel-account-purge.service.ts @@ -131,6 +131,10 @@ export async function purgeArchivedChannelAccount( .deleteFrom("conversations") .where("channel_account_id", "=", accountId) .execute(); + await trx + .deleteFrom("message_reactions") + .where("channel_account_id", "=", accountId) + .execute(); await trx .deleteFrom("channel_accounts") .where("id", "=", accountId)