diff --git a/bun.lock b/bun.lock index db878516e..570713f06 100644 --- a/bun.lock +++ b/bun.lock @@ -211,6 +211,24 @@ "typescript": "catalog:", }, }, + "typescript/discord-pb-dirty-monitor": { + "name": "tachi-discord-pb-dirty-monitor", + "version": "0.1.0", + "dependencies": { + "discord.js": "catalog:", + "dotenv": "catalog:", + "express": "catalog:", + "pg": "catalog:", + }, + "devDependencies": { + "@types/express": "catalog:", + "@types/node": "catalog:", + "@types/pg": "catalog:", + "@typescript/native-preview": "catalog:", + "eslint-config-tachi": "workspace:*", + "typescript": "catalog:", + }, + }, "typescript/docs-autogen-scripts": { "name": "tachi-autogen-docs", "version": "1.0.0", @@ -2894,6 +2912,8 @@ "tachi-db-migration-engine-cli": ["tachi-db-migration-engine-cli@workspace:typescript/db-cli"], + "tachi-discord-pb-dirty-monitor": ["tachi-discord-pb-dirty-monitor@workspace:typescript/discord-pb-dirty-monitor"], + "tachi-github-bot": ["tachi-github-bot@workspace:typescript/github-bot"], "tachi-seeds-scripts": ["tachi-seeds-scripts@workspace:typescript/seeds-scripts"], diff --git a/db/migrations/20260518160000_chart_leaderboard_materialized.sql b/db/migrations/20260518160000_chart_leaderboard_materialized.sql new file mode 100644 index 000000000..f813ccaab --- /dev/null +++ b/db/migrations/20260518160000_chart_leaderboard_materialized.sql @@ -0,0 +1,153 @@ +-- Replace chart_leaderboard VIEW (full-table window agg on every query) with a real table +-- keyed by pb.row_id. Rows are rebuilt per (chart_id, lens) partition when pb changes, +-- using AFTER STATEMENT triggers so bulk INSERT touches each partition at most once. +-- +-- Matches genesis leaderboard ordering (rank tie-breakers + time_achieved ASC NULLS LAST). +-- After migrate: regenerate DB types from your checkout (see db/kanel_config.js). + +CREATE TABLE chart_leaderboard_new ( + row_id UUID PRIMARY KEY REFERENCES pb (row_id) ON DELETE CASCADE, + rank BIGINT NOT NULL, + out_of BIGINT NOT NULL +); + +INSERT INTO chart_leaderboard_new (row_id, rank, out_of) +SELECT + row_id, + RANK() OVER ( + PARTITION BY chart_id, lens + ORDER BY + ranking_value DESC NULLS LAST, + ranking_value_tb1 DESC NULLS LAST, + ranking_value_tb2 DESC NULLS LAST, + ranking_value_tb3 DESC NULLS LAST, + ranking_value_tb4 DESC NULLS LAST, + ranking_value_tb5 DESC NULLS LAST, + time_achieved ASC NULLS LAST + ) AS rank, + COUNT(*) OVER (PARTITION BY chart_id, lens) AS out_of +FROM pb; + +DROP VIEW chart_leaderboard; + +ALTER TABLE chart_leaderboard_new RENAME TO chart_leaderboard; + +COMMENT ON TABLE chart_leaderboard IS + 'Cached global rank / out_of per pb row for PARTITION BY (chart_id, lens). Maintained by triggers on pb; ordering matches genesis chart_leaderboard view.'; + +CREATE FUNCTION refresh_chart_leaderboard_partition ( + p_chart_id TEXT, + p_lens TEXT +) +RETURNS VOID +LANGUAGE plpgsql +SET search_path = public +AS $$ +BEGIN + DELETE FROM chart_leaderboard cl + USING pb + WHERE + cl.row_id = pb.row_id + AND pb.chart_id = p_chart_id + AND pb.lens IS NOT DISTINCT FROM p_lens; + + INSERT INTO chart_leaderboard (row_id, rank, out_of) + SELECT + pb.row_id, + RANK() OVER ( + PARTITION BY pb.chart_id, pb.lens + ORDER BY + pb.ranking_value DESC NULLS LAST, + pb.ranking_value_tb1 DESC NULLS LAST, + pb.ranking_value_tb2 DESC NULLS LAST, + pb.ranking_value_tb3 DESC NULLS LAST, + pb.ranking_value_tb4 DESC NULLS LAST, + pb.ranking_value_tb5 DESC NULLS LAST, + pb.time_achieved ASC NULLS LAST + ) AS rank, + COUNT(*) OVER (PARTITION BY pb.chart_id, pb.lens) AS out_of + FROM pb + WHERE + pb.chart_id = p_chart_id + AND pb.lens IS NOT DISTINCT FROM p_lens; +END; +$$; + +CREATE FUNCTION refresh_leaderboard_stmt_after_pb_insert() +RETURNS TRIGGER +LANGUAGE plpgsql +SET search_path = public +AS $$ +BEGIN + PERFORM refresh_chart_leaderboard_partition (d.chart_id, d.lens) + FROM ( + SELECT DISTINCT + chart_id, + lens + FROM pb_inserted + ) AS d; + + RETURN NULL; +END; +$$; + +CREATE FUNCTION refresh_leaderboard_stmt_after_pb_update() +RETURNS TRIGGER +LANGUAGE plpgsql +SET search_path = public +AS $$ +BEGIN + PERFORM refresh_chart_leaderboard_partition (d.chart_id, d.lens) + FROM ( + SELECT DISTINCT + chart_id, + lens + FROM pb_updated_old + UNION + SELECT DISTINCT + chart_id, + lens + FROM pb_updated_new + ) AS d; + + RETURN NULL; +END; +$$; + +CREATE FUNCTION refresh_leaderboard_stmt_after_pb_delete() +RETURNS TRIGGER +LANGUAGE plpgsql +SET search_path = public +AS $$ +BEGIN + PERFORM refresh_chart_leaderboard_partition (d.chart_id, d.lens) + FROM ( + SELECT DISTINCT + chart_id, + lens + FROM pb_deleted + ) AS d; + + RETURN NULL; +END; +$$; + +CREATE TRIGGER pb_leaderboard_ai + AFTER INSERT ON pb + REFERENCING NEW TABLE AS pb_inserted + FOR EACH STATEMENT + EXECUTE FUNCTION refresh_leaderboard_stmt_after_pb_insert(); + +CREATE TRIGGER pb_leaderboard_au + AFTER UPDATE ON pb + REFERENCING OLD TABLE AS pb_updated_old NEW TABLE AS pb_updated_new + FOR EACH STATEMENT + EXECUTE FUNCTION refresh_leaderboard_stmt_after_pb_update(); + +CREATE TRIGGER pb_leaderboard_ad + AFTER DELETE ON pb + REFERENCING OLD TABLE AS pb_deleted + FOR EACH STATEMENT + EXECUTE FUNCTION refresh_leaderboard_stmt_after_pb_delete(); + +ANALYZE chart_leaderboard; diff --git a/typescript/db/src/generated/index.ts b/typescript/db/src/generated/index.ts index d7cd776c3..e3a505d6a 100644 --- a/typescript/db/src/generated/index.ts +++ b/typescript/db/src/generated/index.ts @@ -50,6 +50,7 @@ export { type default as QuestSubTable, type QuestSub, type NewQuestSub, type Qu export { type orphan_score_row_id, type default as OrphanScoreTable, type OrphanScore, type NewOrphanScore, type OrphanScoreUpdate } from './public/OrphanScore'; export { type default as FolderViewTable, type FolderView, type NewFolderView, type FolderViewUpdate } from './public/FolderView'; export { type table_id, type default as TableTable, type Table, type NewTable, type TableUpdate } from './public/Table'; +export { type default as ChartLeaderboardTable, type ChartLeaderboard, type NewChartLeaderboard, type ChartLeaderboardUpdate } from './public/ChartLeaderboard'; export { type default as GoalSubTable, type GoalSub, type NewGoalSub, type GoalSubUpdate } from './public/GoalSub'; export { type job_queue_row_id, type default as JobQueueTable, type JobQueue, type NewJobQueue, type JobQueueUpdate } from './public/JobQueue'; export { type priv_api_client_client_id, type default as PrivApiClientTable, type PrivApiClient, type NewPrivApiClient, type PrivApiClientUpdate } from './public/PrivApiClient'; @@ -67,7 +68,6 @@ export { type questline_id, type default as QuestlineTable, type Questline, type export { type cron_task_id, type default as CronTaskTable, type CronTask, type NewCronTask, type CronTaskUpdate } from './public/CronTask'; export { type import_class_row_id, type default as ImportClassTable, type ImportClass, type NewImportClass, type ImportClassUpdate } from './public/ImportClass'; export { type game_stats_snapshot_game, type game_stats_snapshot_timestamp, type default as GameStatsSnapshotTable, type GameStatsSnapshot, type NewGameStatsSnapshot, type GameStatsSnapshotUpdate } from './public/GameStatsSnapshot'; -export { type default as ChartLeaderboardTable, type ChartLeaderboard } from './public/ChartLeaderboard'; export { type default as AuthLevel } from './public/AuthLevel'; export { type default as AccountBadgeKind } from './public/AccountBadgeKind'; export { type default as ImportStatus } from './public/ImportStatus'; @@ -75,9 +75,13 @@ export { type default as Game } from './public/Game'; export { type default as ActionResult } from './public/ActionResult'; export { type default as GameGroup } from './public/GameGroup'; export { type default as ImportType } from './public/ImportType'; +export { type refresh_leaderboard_stmt_after_pb_insert_params } from './public/refresh_leaderboard_stmt_after_pb_insert'; export { type enqueue_pb_dirty_params } from './public/enqueue_pb_dirty'; export { type enqueue_game_profile_dirty_params } from './public/enqueue_game_profile_dirty'; export { type enqueue_score_rederive_params } from './public/enqueue_score_rederive'; +export { type refresh_leaderboard_stmt_after_pb_delete_params } from './public/refresh_leaderboard_stmt_after_pb_delete'; export { type enqueue_session_dirty_params } from './public/enqueue_session_dirty'; +export { type refresh_leaderboard_stmt_after_pb_update_params } from './public/refresh_leaderboard_stmt_after_pb_update'; +export { type refresh_chart_leaderboard_partition_params } from './public/refresh_chart_leaderboard_partition'; export { type default as PublicSchema } from './public/PublicSchema'; export { type default as Database } from './Database'; diff --git a/typescript/db/src/generated/public/ChartLeaderboard.ts b/typescript/db/src/generated/public/ChartLeaderboard.ts index ce2ac4f61..e216edaa4 100644 --- a/typescript/db/src/generated/public/ChartLeaderboard.ts +++ b/typescript/db/src/generated/public/ChartLeaderboard.ts @@ -2,47 +2,22 @@ // This file is automatically generated by Kanel. Do not modify manually. import type { pb_row_id } from './Pb'; -import type { account_id } from './Account'; -import type { chart_id } from './Chart'; -import type { ColumnType, Selectable } from 'kysely'; +import type { ColumnType, Selectable, Insertable, Updateable } from 'kysely'; -/** Represents the view public.chart_leaderboard */ +/** + * Represents the table public.chart_leaderboard + * Cached global rank / out_of per pb row for PARTITION BY (chart_id, lens). Maintained by triggers on pb; ordering matches genesis chart_leaderboard view. + */ export default interface ChartLeaderboardTable { - row_id: ColumnType; + row_id: ColumnType; - user_id: ColumnType; + rank: ColumnType; - chart_id: ColumnType; - - lens: ColumnType; - - data: ColumnType; - - derived_data: ColumnType; - - calculated_data: ColumnType; - - judgements: ColumnType; - - ranking_value: ColumnType; - - ranking_value_tb1: ColumnType; - - ranking_value_tb2: ColumnType; - - ranking_value_tb3: ColumnType; - - ranking_value_tb4: ColumnType; - - ranking_value_tb5: ColumnType; - - highlight: ColumnType; - - time_achieved: ColumnType; - - rank: ColumnType; - - out_of: ColumnType; + out_of: ColumnType; } export type ChartLeaderboard = Selectable; + +export type NewChartLeaderboard = Insertable; + +export type ChartLeaderboardUpdate = Updateable; diff --git a/typescript/db/src/generated/public/PublicSchema.ts b/typescript/db/src/generated/public/PublicSchema.ts index 9aa8fa958..8d8b02cb9 100644 --- a/typescript/db/src/generated/public/PublicSchema.ts +++ b/typescript/db/src/generated/public/PublicSchema.ts @@ -50,6 +50,7 @@ import type { default as QuestSubTable } from './QuestSub'; import type { default as OrphanScoreTable } from './OrphanScore'; import type { default as FolderViewTable } from './FolderView'; import type { default as TableTable } from './Table'; +import type { default as ChartLeaderboardTable } from './ChartLeaderboard'; import type { default as GoalSubTable } from './GoalSub'; import type { default as JobQueueTable } from './JobQueue'; import type { default as PrivApiClientTable } from './PrivApiClient'; @@ -67,7 +68,6 @@ import type { default as QuestlineTable } from './Questline'; import type { default as CronTaskTable } from './CronTask'; import type { default as ImportClassTable } from './ImportClass'; import type { default as GameStatsSnapshotTable } from './GameStatsSnapshot'; -import type { default as ChartLeaderboardTable } from './ChartLeaderboard'; export default interface PublicSchema { priv_svc_cg_card_info: PrivSvcCgCardInfoTable; @@ -168,6 +168,8 @@ export default interface PublicSchema { table: TableTable; + chart_leaderboard: ChartLeaderboardTable; + goal_sub: GoalSubTable; job_queue: JobQueueTable; @@ -201,6 +203,4 @@ export default interface PublicSchema { import_class: ImportClassTable; game_stats_snapshot: GameStatsSnapshotTable; - - chart_leaderboard: ChartLeaderboardTable; } diff --git a/typescript/db/src/generated/public/refresh_chart_leaderboard_partition.ts b/typescript/db/src/generated/public/refresh_chart_leaderboard_partition.ts new file mode 100644 index 000000000..c53c6edc6 --- /dev/null +++ b/typescript/db/src/generated/public/refresh_chart_leaderboard_partition.ts @@ -0,0 +1,8 @@ +// @generated +// This file is automatically generated by Kanel. Do not modify manually. + +export interface refresh_chart_leaderboard_partition_params { + p_chart_id: string; + + p_lens: string; +} diff --git a/typescript/db/src/generated/public/refresh_leaderboard_stmt_after_pb_delete.ts b/typescript/db/src/generated/public/refresh_leaderboard_stmt_after_pb_delete.ts new file mode 100644 index 000000000..7f08bc046 --- /dev/null +++ b/typescript/db/src/generated/public/refresh_leaderboard_stmt_after_pb_delete.ts @@ -0,0 +1,5 @@ +// @generated +// This file is automatically generated by Kanel. Do not modify manually. + +export interface refresh_leaderboard_stmt_after_pb_delete_params { +} diff --git a/typescript/db/src/generated/public/refresh_leaderboard_stmt_after_pb_insert.ts b/typescript/db/src/generated/public/refresh_leaderboard_stmt_after_pb_insert.ts new file mode 100644 index 000000000..97c74fc7f --- /dev/null +++ b/typescript/db/src/generated/public/refresh_leaderboard_stmt_after_pb_insert.ts @@ -0,0 +1,5 @@ +// @generated +// This file is automatically generated by Kanel. Do not modify manually. + +export interface refresh_leaderboard_stmt_after_pb_insert_params { +} diff --git a/typescript/db/src/generated/public/refresh_leaderboard_stmt_after_pb_update.ts b/typescript/db/src/generated/public/refresh_leaderboard_stmt_after_pb_update.ts new file mode 100644 index 000000000..45d8729a1 --- /dev/null +++ b/typescript/db/src/generated/public/refresh_leaderboard_stmt_after_pb_update.ts @@ -0,0 +1,5 @@ +// @generated +// This file is automatically generated by Kanel. Do not modify manually. + +export interface refresh_leaderboard_stmt_after_pb_update_params { +} diff --git a/typescript/eslint-config/rules/require-unicode-regexp-fix.js b/typescript/eslint-config/rules/require-unicode-regexp-fix.js index 50e4ddb5b..79028f1a1 100644 --- a/typescript/eslint-config/rules/require-unicode-regexp-fix.js +++ b/typescript/eslint-config/rules/require-unicode-regexp-fix.js @@ -134,8 +134,9 @@ export const requireUnicodeRegexpFix = { context.report({ messageId: - requireFlag === "v" /** @type {"requireVFlag"|"requireUFlag"} */ - ? "requireVFlag" + requireFlag === "v" + ? /** @type {"requireVFlag"|"requireUFlag"} */ + "requireVFlag" : "requireUFlag", node, fix: isValidWithUnicodeFlag( @@ -149,7 +150,9 @@ export const requireUnicodeRegexpFix = { }, Program(node) { - const scope = /** @type {import("eslint").Scope.Scope} */ (sourceCode.getScope(node)); + const scope = /** @type {import("eslint").Scope.Scope} */ ( + sourceCode.getScope(node) + ); const tracker = new ReferenceTracker(scope); const trackMap = { RegExp: { [CALL]: true, [CONSTRUCT]: true } }; @@ -157,17 +160,12 @@ export const requireUnicodeRegexpFix = { for (const { node: refNode } of tracker.iterateGlobalReferences(trackMap)) { const [patternNode, flagsNode] = refNode.arguments; - if ( - patternNode && - patternNode.type === "SpreadElement" - ) { + if (patternNode && patternNode.type === "SpreadElement") { continue; } - const pattern = - getStringIfConstant(patternNode, scope); - const flags = - flagsNode ? getStringIfConstant(flagsNode, scope) : undefined; + const pattern = getStringIfConstant(patternNode, scope); + const flags = flagsNode ? getStringIfConstant(flagsNode, scope) : undefined; let missingFlag = !flagsNode; @@ -189,22 +187,22 @@ export const requireUnicodeRegexpFix = { context.report({ messageId: - requireFlag === "v" /** @type {"requireVFlag"|"requireUFlag"} */ - ? "requireVFlag" + requireFlag === "v" + ? /** @type {"requireVFlag"|"requireUFlag"} */ + "requireVFlag" : "requireUFlag", node: refNode, - fix: - canFixPattern - ? (fixer) => - fixGlobalRegExpCall( - fixer, - sourceCode, - refNode, - flagsNode, - flags ?? "", - requireFlag, - ) - : null, + fix: canFixPattern + ? (fixer) => + fixGlobalRegExpCall( + fixer, + sourceCode, + refNode, + flagsNode, + flags ?? "", + requireFlag, + ) + : null, }); } }, @@ -225,12 +223,7 @@ function fixRegexLiteral(fixer, sourceCode, node, requireFlag) { if (requireFlag) { const conflicting = requireFlag === "u" /** @type {"u"|"v"} */ ? "v" : "u"; - if ( - regexText.includes( - conflicting, - slashPos, - ) - ) { + if (regexText.includes(conflicting, slashPos)) { return fixer.replaceText( node, regexText.slice(0, slashPos) + @@ -252,14 +245,7 @@ function fixRegexLiteral(fixer, sourceCode, node, requireFlag) { * @param {string} flags Resolved constant flags when determinable (`""` if omitted) * @param {RequireFlagOption} requireFlag */ -function fixGlobalRegExpCall( - fixer, - sourceCode, - refNode, - flagsNode, - flags, - requireFlag, -) { +function fixGlobalRegExpCall(fixer, sourceCode, refNode, flagsNode, flags, requireFlag) { const replaceFlag = requireFlag ?? /** @type {"u"|"v"} */ @@ -284,10 +270,7 @@ function fixGlobalRegExpCall( if ( flagsNode.type === "TemplateLiteral" && (flagsNode.expressions.length || - flagsNode.quasis.some( - (q) => - q.value.raw.includes("\\"), - )) + flagsNode.quasis.some((q) => q.value.raw.includes("\\"))) ) { return null; } @@ -298,28 +281,25 @@ function fixGlobalRegExpCall( ); } - return fixer.replaceText(flagsNode, [ - flagsNodeText.slice(0, flagsNodeText.length - 1), - flagsNodeText.slice(flagsNodeText.length - 1), - ].join(replaceFlag)); + return fixer.replaceText( + flagsNode, + [ + flagsNodeText.slice(0, flagsNodeText.length - 1), + flagsNodeText.slice(flagsNodeText.length - 1), + ].join(replaceFlag), + ); } return null; } - const penultimateToken = - sourceCode.getLastToken( - refNode, - { skip: 1 }, - ); + const penultimateToken = sourceCode.getLastToken(refNode, { skip: 1 }); if (!penultimateToken) { return null; } return fixer.insertTextAfter( penultimateToken, - isCommaToken(penultimateToken) - ? ` "${replaceFlag}",` - : `, "${replaceFlag}"`, + isCommaToken(penultimateToken) ? ` "${replaceFlag}",` : `, "${replaceFlag}"`, ); } diff --git a/typescript/server/src/lib/jobs/cron/cron-service.ts b/typescript/server/src/lib/jobs/cron/cron-service.ts index 84268c34a..94bec3990 100644 --- a/typescript/server/src/lib/jobs/cron/cron-service.ts +++ b/typescript/server/src/lib/jobs/cron/cron-service.ts @@ -1,8 +1,10 @@ +import type { Database } from "tachi-db"; + import { getCronTaskDefinitions } from "#lib/jobs/cron/cron-registry"; import { log } from "#lib/log/log"; import DB from "#services/pg/db"; import { CronExpressionParser } from "cron-parser"; -import { sql } from "kysely"; +import { type Kysely, sql } from "kysely"; /** Namespaces `pg_try_advisory_lock` for the cron scheduler. */ const CRON_ADVISORY_KEY1 = 0x54_61_63_68; // "Tach" @@ -65,11 +67,12 @@ export function getDueFireTime(schedule: string, last: Date | null, now: Date): throw new Error(`getDueFireTime exceeded ${maxSteps} next() steps (schedule ${schedule}).`); } -export async function syncCronTasksFromRegistry(): Promise { +export async function syncCronTasksFromRegistry(executor: Kysely = DB): Promise { const defs = getCronTaskDefinitions(); const now = new Date().toISOString(); for (const d of defs) { - await DB.insertInto("cron_task") + await executor + .insertInto("cron_task") .values({ id: d.id, schedule: d.schedule, @@ -89,93 +92,97 @@ export async function syncCronTasksFromRegistry(): Promise { } } -async function tryAcquireCronTickLock(): Promise { - const r = await sql<{ acquired: boolean }>` - SELECT pg_try_advisory_lock(${CRON_ADVISORY_KEY1}, ${CRON_ADVISORY_KEY2}) AS acquired - `.execute(DB); - const row = r.rows[0] as { acquired: boolean } | undefined; - return row?.acquired === true; -} - -async function releaseCronTickLock(): Promise { - await sql`SELECT pg_advisory_unlock(${CRON_ADVISORY_KEY1}, ${CRON_ADVISORY_KEY2})`.execute(DB); -} - +/** + * Session advisory locks must be released on the same backend that acquired them. The global + * pool used by `DB` would otherwise acquire on connection A and `pg_advisory_unlock` on B, + * leaving the lock held until A disconnects — then every `pg_try_advisory_lock` fails and + * `runCronTickOnce` no-ops with no INFO-level log. + */ export async function runCronTickOnce(): Promise { - const got = await tryAcquireCronTickLock(); - if (!got) { - return; - } - try { - await syncCronTasksFromRegistry(); - const now = new Date(); - const rows = await DB.selectFrom("cron_task") - .select(["cron_task.id", "cron_task.schedule", "cron_task.last_scheduled_at"]) - .execute(); - const byId = new Map(rows.map((r) => [r.id, r])); - for (const def of getCronTaskDefinitions()) { - const row = byId.get(def.id); - if (!row) { - continue; - } - const last = row.last_scheduled_at ? new Date(row.last_scheduled_at) : null; - const due = getDueFireTime(def.schedule, last, now); - if (!due) { - continue; - } - const scheduledAtIso = due.toISOString(); - const execRow = await DB.insertInto("cron_task_execution") - .values({ - task_id: def.id, - scheduled_at: scheduledAtIso, - status: "running", - completed_at: null, - output: null, - error: null, - }) - .returning("cron_task_execution.id") - .executeTakeFirstOrThrow(); - - try { - await def.run(); - await DB.updateTable("cron_task") - .set({ - last_scheduled_at: scheduledAtIso, - updated_at: new Date().toISOString(), - }) - .where("cron_task.id", "=", def.id) - .execute(); - await DB.updateTable("cron_task_execution") - .set({ - status: "success", - completed_at: new Date().toISOString(), + await DB.connection().execute(async (conn) => { + const r = await sql<{ acquired: boolean }>` + SELECT pg_try_advisory_lock(${CRON_ADVISORY_KEY1}, ${CRON_ADVISORY_KEY2}) AS acquired + `.execute(conn); + const got = (r.rows[0] as { acquired: boolean } | undefined)?.acquired === true; + if (!got) { + log.debug( + "Cron tick skipped: advisory lock not acquired (another worker or stale session lock).", + ); + return; + } + try { + await syncCronTasksFromRegistry(conn); + const now = new Date(); + const rows = await DB.selectFrom("cron_task") + .select(["cron_task.id", "cron_task.schedule", "cron_task.last_scheduled_at"]) + .execute(); + const byId = new Map(rows.map((r) => [r.id, r])); + for (const def of getCronTaskDefinitions()) { + const row = byId.get(def.id); + if (!row) { + continue; + } + const last = row.last_scheduled_at ? new Date(row.last_scheduled_at) : null; + const due = getDueFireTime(def.schedule, last, now); + if (!due) { + continue; + } + const scheduledAtIso = due.toISOString(); + const execRow = await DB.insertInto("cron_task_execution") + .values({ + task_id: def.id, + scheduled_at: scheduledAtIso, + status: "running", + completed_at: null, output: null, error: null, }) - .where("cron_task_execution.id", "=", execRow.id) - .execute(); - log.info(`Cron task ${def.id} completed for fire ${scheduledAtIso}.`); - } catch (err) { - const message = err instanceof Error ? err.message : String(err); - log.error({ err }, `Cron task ${def.id} failed.`); - await DB.updateTable("cron_task") - .set({ - last_scheduled_at: scheduledAtIso, - updated_at: new Date().toISOString(), - }) - .where("cron_task.id", "=", def.id) - .execute(); - await DB.updateTable("cron_task_execution") - .set({ - status: "failure", - completed_at: new Date().toISOString(), - error: message, - }) - .where("cron_task_execution.id", "=", execRow.id) - .execute(); + .returning("cron_task_execution.id") + .executeTakeFirstOrThrow(); + + try { + await def.run(); + await DB.updateTable("cron_task") + .set({ + last_scheduled_at: scheduledAtIso, + updated_at: new Date().toISOString(), + }) + .where("cron_task.id", "=", def.id) + .execute(); + await DB.updateTable("cron_task_execution") + .set({ + status: "success", + completed_at: new Date().toISOString(), + output: null, + error: null, + }) + .where("cron_task_execution.id", "=", execRow.id) + .execute(); + log.info(`Cron task ${def.id} completed for fire ${scheduledAtIso}.`); + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + log.error({ err }, `Cron task ${def.id} failed.`); + await DB.updateTable("cron_task") + .set({ + last_scheduled_at: scheduledAtIso, + updated_at: new Date().toISOString(), + }) + .where("cron_task.id", "=", def.id) + .execute(); + await DB.updateTable("cron_task_execution") + .set({ + status: "failure", + completed_at: new Date().toISOString(), + error: message, + }) + .where("cron_task_execution.id", "=", execRow.id) + .execute(); + } } + } finally { + await sql` + SELECT pg_advisory_unlock(${CRON_ADVISORY_KEY1}, ${CRON_ADVISORY_KEY2}) + `.execute(conn); } - } finally { - await releaseCronTickLock(); - } + }); }