From 218408493ae908c26c8573f1a695fd6663b35b7e Mon Sep 17 00:00:00 2001 From: x032205 Date: Fri, 11 Jul 2025 22:05:32 -0400 Subject: [PATCH 1/3] Optimize token cleanup job --- .../identity-access-token-dal.ts | 43 ++++++++----------- 1 file changed, 19 insertions(+), 24 deletions(-) diff --git a/backend/src/services/identity-access-token/identity-access-token-dal.ts b/backend/src/services/identity-access-token/identity-access-token-dal.ts index 8f65d8555..3b6b412ba 100644 --- a/backend/src/services/identity-access-token/identity-access-token-dal.ts +++ b/backend/src/services/identity-access-token/identity-access-token-dal.ts @@ -31,9 +31,10 @@ export const identityAccessTokenDALFactory = (db: TDbClient) => { logger.info(`${QueueName.DailyResourceCleanUp}: remove expired access token started`); const MAX_TTL = 315_360_000; // Maximum TTL value in seconds (10 years) + const QUERY_TIMEOUT_MS = 10 * 60 * 1000; // 10 minutes - try { - const docs = (tx || db)(TableName.IdentityAccessToken) + const performDelete = (dbClient: Knex | Knex.Transaction) => + dbClient(TableName.IdentityAccessToken) .where({ isAccessTokenRevoked: true }) @@ -47,30 +48,24 @@ export const identityAccessTokenDALFactory = (db: TDbClient) => { ); }) .orWhere((qb) => { - void qb.where("accessTokenTTL", ">", 0).andWhere((qb2) => { - void qb2 - .where((qb3) => { - void qb3 - .whereNotNull("accessTokenLastRenewedAt") - // accessTokenLastRenewedAt + convert_integer_to_seconds(accessTokenTTL) < present_date - .andWhereRaw( - `"${TableName.IdentityAccessToken}"."accessTokenLastRenewedAt" + make_interval(secs => LEAST("${TableName.IdentityAccessToken}"."accessTokenTTL", ?)) < NOW()`, - [MAX_TTL] - ); - }) - .orWhere((qb3) => { - void qb3 - .whereNull("accessTokenLastRenewedAt") - // created + convert_integer_to_seconds(accessTokenTTL) < present_date - .andWhereRaw( - `"${TableName.IdentityAccessToken}"."createdAt" + make_interval(secs => LEAST("${TableName.IdentityAccessToken}"."accessTokenTTL", ?)) < NOW()`, - [MAX_TTL] - ); - }); - }); + void qb + .where("accessTokenTTL", ">", 0) + .andWhereRaw( + `COALESCE("${TableName.IdentityAccessToken}"."accessTokenLastRenewedAt", "${TableName.IdentityAccessToken}"."createdAt") + make_interval(secs => LEAST("${TableName.IdentityAccessToken}"."accessTokenTTL", ?)) < NOW()`, + [MAX_TTL] + ); }) .delete(); - await docs; + + try { + if (tx) { + await performDelete(tx); + } else { + await db.transaction(async (trx) => { + await trx.raw(`SET statement_timeout = ${QUERY_TIMEOUT_MS}`); + await performDelete(trx); + }); + } logger.info(`${QueueName.DailyResourceCleanUp}: remove expired access token completed`); } catch (error) { throw new DatabaseError({ error, name: "IdentityAccessTokenPrune" }); From 513f942aae5966b208b5242fa284b6056409c771 Mon Sep 17 00:00:00 2001 From: x032205 Date: Mon, 14 Jul 2025 00:39:34 -0400 Subject: [PATCH 2/3] Add batching to not lock DB --- .../identity-access-token-dal.ts | 58 ++++++++++++++----- 1 file changed, 44 insertions(+), 14 deletions(-) diff --git a/backend/src/services/identity-access-token/identity-access-token-dal.ts b/backend/src/services/identity-access-token/identity-access-token-dal.ts index 3b6b412ba..b6b8ad425 100644 --- a/backend/src/services/identity-access-token/identity-access-token-dal.ts +++ b/backend/src/services/identity-access-token/identity-access-token-dal.ts @@ -30,10 +30,16 @@ export const identityAccessTokenDALFactory = (db: TDbClient) => { const removeExpiredTokens = async (tx?: Knex) => { logger.info(`${QueueName.DailyResourceCleanUp}: remove expired access token started`); - const MAX_TTL = 315_360_000; // Maximum TTL value in seconds (10 years) + const BATCH_SIZE = 10000; + const MAX_RETRY_ON_FAILURE = 3; const QUERY_TIMEOUT_MS = 10 * 60 * 1000; // 10 minutes + const MAX_TTL = 315_360_000; // Maximum TTL value in seconds (10 years) - const performDelete = (dbClient: Knex | Knex.Transaction) => + let deletedTokenIds: { id: string }[] = []; + let numberOfRetryOnFailure = 0; + let isRetrying = false; + + const getExpiredTokensQuery = (dbClient: Knex | Knex.Transaction) => dbClient(TableName.IdentityAccessToken) .where({ isAccessTokenRevoked: true @@ -54,22 +60,46 @@ export const identityAccessTokenDALFactory = (db: TDbClient) => { `COALESCE("${TableName.IdentityAccessToken}"."accessTokenLastRenewedAt", "${TableName.IdentityAccessToken}"."createdAt") + make_interval(secs => LEAST("${TableName.IdentityAccessToken}"."accessTokenTTL", ?)) < NOW()`, [MAX_TTL] ); - }) - .delete(); + }); - try { - if (tx) { - await performDelete(tx); - } else { - await db.transaction(async (trx) => { - await trx.raw(`SET statement_timeout = ${QUERY_TIMEOUT_MS}`); - await performDelete(trx); + do { + try { + const deleteBatch = async (dbClient: Knex | Knex.Transaction) => { + const idsToDeleteQuery = getExpiredTokensQuery(dbClient).select("id").limit(BATCH_SIZE); + return dbClient(TableName.IdentityAccessToken).whereIn("id", idsToDeleteQuery).del().returning("id"); + }; + + if (tx) { + // eslint-disable-next-line no-await-in-loop + deletedTokenIds = await deleteBatch(tx); + } else { + // eslint-disable-next-line no-await-in-loop + deletedTokenIds = await db.transaction(async (trx) => { + await trx.raw(`SET statement_timeout = ${QUERY_TIMEOUT_MS}`); + return deleteBatch(trx); + }); + } + + numberOfRetryOnFailure = 0; // reset + } catch (error) { + numberOfRetryOnFailure += 1; + logger.error(error, "Failed to delete a batch of expired identity access tokens on pruning"); + } finally { + // eslint-disable-next-line no-await-in-loop + await new Promise((resolve) => { + setTimeout(resolve, 10); // time to breathe for db }); } - logger.info(`${QueueName.DailyResourceCleanUp}: remove expired access token completed`); - } catch (error) { - throw new DatabaseError({ error, name: "IdentityAccessTokenPrune" }); + isRetrying = numberOfRetryOnFailure > 0; + } while (deletedTokenIds.length > 0 || (isRetrying && numberOfRetryOnFailure < MAX_RETRY_ON_FAILURE)); + + if (numberOfRetryOnFailure >= MAX_RETRY_ON_FAILURE) { + logger.error( + `IdentityAccessTokenPrune: Pruning failed and stopped after ${MAX_RETRY_ON_FAILURE} consecutive retries.` + ); } + + logger.info(`${QueueName.DailyResourceCleanUp}: remove expired access token completed`); }; return { ...identityAccessTokenOrm, findOne, removeExpiredTokens }; From 2f0a247c115532f427c982ba211cbc53d209a621 Mon Sep 17 00:00:00 2001 From: x032205 Date: Mon, 14 Jul 2025 18:01:35 -0400 Subject: [PATCH 3/3] Describe query --- .../identity-access-token-dal.ts | 24 ++++++++++++++----- 1 file changed, 18 insertions(+), 6 deletions(-) diff --git a/backend/src/services/identity-access-token/identity-access-token-dal.ts b/backend/src/services/identity-access-token/identity-access-token-dal.ts index b6b8ad425..19de362d8 100644 --- a/backend/src/services/identity-access-token/identity-access-token-dal.ts +++ b/backend/src/services/identity-access-token/identity-access-token-dal.ts @@ -54,12 +54,24 @@ export const identityAccessTokenDALFactory = (db: TDbClient) => { ); }) .orWhere((qb) => { - void qb - .where("accessTokenTTL", ">", 0) - .andWhereRaw( - `COALESCE("${TableName.IdentityAccessToken}"."accessTokenLastRenewedAt", "${TableName.IdentityAccessToken}"."createdAt") + make_interval(secs => LEAST("${TableName.IdentityAccessToken}"."accessTokenTTL", ?)) < NOW()`, - [MAX_TTL] - ); + void qb.where("accessTokenTTL", ">", 0).andWhereRaw( + ` + -- Check if the token's effective expiration time has passed. + -- The expiration time is calculated by adding its TTL to its last renewal/creation time. + COALESCE( + "${TableName.IdentityAccessToken}"."accessTokenLastRenewedAt", -- Use last renewal time if available + "${TableName.IdentityAccessToken}"."createdAt" -- Otherwise, use creation time + ) + + make_interval( + secs => LEAST( + "${TableName.IdentityAccessToken}"."accessTokenTTL", -- Token's specified TTL + ? -- Capped by MAX_TTL (parameterized value) + ) + ) + < NOW() -- Check if the calculated time is before now + `, + [MAX_TTL] + ); }); do {