mirror of
https://github.com/zkldi/Tachi.git
synced 2026-10-05 13:28:08 +03:00
fix: db optimisations, 2
This commit is contained in:
@@ -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"],
|
||||
|
||||
@@ -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;
|
||||
@@ -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';
|
||||
|
||||
@@ -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<pb_row_id, never, never>;
|
||||
row_id: ColumnType<pb_row_id, pb_row_id, pb_row_id>;
|
||||
|
||||
user_id: ColumnType<account_id, never, never>;
|
||||
rank: ColumnType<number, number, number>;
|
||||
|
||||
chart_id: ColumnType<chart_id, never, never>;
|
||||
|
||||
lens: ColumnType<string | null, never, never>;
|
||||
|
||||
data: ColumnType<unknown, never, never>;
|
||||
|
||||
derived_data: ColumnType<unknown, never, never>;
|
||||
|
||||
calculated_data: ColumnType<unknown, never, never>;
|
||||
|
||||
judgements: ColumnType<unknown, never, never>;
|
||||
|
||||
ranking_value: ColumnType<number, never, never>;
|
||||
|
||||
ranking_value_tb1: ColumnType<number | null, never, never>;
|
||||
|
||||
ranking_value_tb2: ColumnType<number | null, never, never>;
|
||||
|
||||
ranking_value_tb3: ColumnType<number | null, never, never>;
|
||||
|
||||
ranking_value_tb4: ColumnType<number | null, never, never>;
|
||||
|
||||
ranking_value_tb5: ColumnType<number | null, never, never>;
|
||||
|
||||
highlight: ColumnType<boolean, never, never>;
|
||||
|
||||
time_achieved: ColumnType<string | null, never, never>;
|
||||
|
||||
rank: ColumnType<number, never, never>;
|
||||
|
||||
out_of: ColumnType<number, never, never>;
|
||||
out_of: ColumnType<number, number, number>;
|
||||
}
|
||||
|
||||
export type ChartLeaderboard = Selectable<ChartLeaderboardTable>;
|
||||
|
||||
export type NewChartLeaderboard = Insertable<ChartLeaderboardTable>;
|
||||
|
||||
export type ChartLeaderboardUpdate = Updateable<ChartLeaderboardTable>;
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
@@ -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 {
|
||||
}
|
||||
@@ -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 {
|
||||
}
|
||||
@@ -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 {
|
||||
}
|
||||
@@ -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}"`,
|
||||
);
|
||||
}
|
||||
|
||||
@@ -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<void> {
|
||||
export async function syncCronTasksFromRegistry(executor: Kysely<Database> = DB): Promise<void> {
|
||||
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<void> {
|
||||
}
|
||||
}
|
||||
|
||||
async function tryAcquireCronTickLock(): Promise<boolean> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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();
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user