diff --git a/.github/workflows/bot.yml b/.github/workflows/bot.yml index 93b1d3637..e9faa11b8 100644 --- a/.github/workflows/bot.yml +++ b/.github/workflows/bot.yml @@ -44,7 +44,7 @@ jobs: NODE_ENV: "test" services: tachi-postgres-test: - image: postgres:18 + image: ghcr.io/zkldi/zk-postgres:18@sha256:484fa9e6501ec93c5729ba9f9312d8ee4dfe0c136f1b4881fab5e0aed5aa68fd env: POSTGRES_USER: tachi POSTGRES_PASSWORD: tachi diff --git a/.github/workflows/server.yml b/.github/workflows/server.yml index ded9893dc..9093dd8f4 100644 --- a/.github/workflows/server.yml +++ b/.github/workflows/server.yml @@ -51,7 +51,7 @@ jobs: # the postgres defaults because `services:` does not let us pass # postgres CLI args. With tmpfs, fsync is essentially free anyway. tachi-postgres-test: - image: postgres:18 + image: ghcr.io/zkldi/zk-postgres:18@sha256:484fa9e6501ec93c5729ba9f9312d8ee4dfe0c136f1b4881fab5e0aed5aa68fd env: POSTGRES_USER: tachi POSTGRES_PASSWORD: tachi diff --git a/db/migrations/20260618150000_score_view.sql b/db/migrations/20260618150000_score_view.sql new file mode 100644 index 000000000..4d71bf177 --- /dev/null +++ b/db/migrations/20260618150000_score_view.sql @@ -0,0 +1,11 @@ +-- Replace the `score` table with a view that only surfaces committed rows. +-- This ensures all public API queries automatically see committed-only scores, +-- while import-internal code can write/read staged rows via `raw_score` directly. +-- +-- Postgres carries all indexes, constraints, triggers, and FK references along +-- with the renamed table automatically (they track by OID, not name). + +ALTER TABLE score RENAME TO raw_score; + +CREATE VIEW score AS + SELECT * FROM raw_score WHERE committed = true; diff --git a/db/migrations/20260618160000_action_partition.sql b/db/migrations/20260618160000_action_partition.sql new file mode 100644 index 000000000..a563cce51 --- /dev/null +++ b/db/migrations/20260618160000_action_partition.sql @@ -0,0 +1,70 @@ +-- Convert `action` to a monthly range-partitioned table managed by pg_partman. +-- +-- Strategy: rename the existing unpartitioned table, recreate it as PARTITION BY RANGE, +-- let pg_partman build the child partitions, copy the existing rows in, then drop the +-- old table. Safe to do in one migration because `action` is small (~120 k rows in +-- prod, no inbound foreign keys) and all reads/writes are blocked during the migration +-- anyway by the lock on the renamed table. +-- +-- pg_partman must be installed on the server before this migration runs: +-- apt-get install postgresql-18-partman +-- CREATE EXTENSION pg_partman SCHEMA partman; -- or let postgres-init.sql do it + +CREATE SCHEMA IF NOT EXISTS partman; +CREATE EXTENSION IF NOT EXISTS pg_partman SCHEMA partman; + +ALTER TABLE action RENAME TO action_unpartitioned; + +CREATE TABLE action ( + row_id UUID DEFAULT uuidv7() NOT NULL, + user_id BIGINT REFERENCES account(id), + ip INET, + app TEXT NOT NULL, + kind TEXT NOT NULL, + result ACTION_RESULT NOT NULL, + input JSONB NOT NULL, + output JSONB, + ts_start TIMESTAMPTZ NOT NULL, + ts_end TIMESTAMPTZ NOT NULL, + -- Partition key must be part of the PK in Postgres partitioned tables. + PRIMARY KEY (row_id, ts_start) +) PARTITION BY RANGE (ts_start); + +-- Indexes on the parent are inherited by every child partition pg_partman creates. +CREATE INDEX ON action (user_id, ts_start DESC); +CREATE INDEX ON action (ip, ts_start DESC); +CREATE INDEX ON action (app, kind, ts_start DESC); + +SELECT partman.create_parent( + p_parent_table => 'public.action', + p_control => 'ts_start', + p_interval => '1 month', + p_premake => 3, + p_start_partition => ( + SELECT COALESCE(MIN(ts_start), NOW())::date::text + FROM action_unpartitioned + ) +); + +-- Keep 24 months of action logs; older partitions are dropped on maintenance runs. +UPDATE partman.part_config +SET retention = '24 months', + retention_keep_table = false +WHERE parent_table = 'public.action'; + +-- Migrate existing rows. +INSERT INTO action SELECT * FROM action_unpartitioned; + +DROP TABLE action_unpartitioned; + +-- Redistribute any rows that landed in the DEFAULT partition into proper monthly slices. +DO $$ +DECLARE + moved BIGINT; +BEGIN + LOOP + SELECT partman.partition_data_time('public.action', p_batch_count => 1000) + INTO moved; + EXIT WHEN moved = 0; + END LOOP; +END $$; diff --git a/db/migrations/20260618170000_game_stats_snapshot_partition.sql b/db/migrations/20260618170000_game_stats_snapshot_partition.sql new file mode 100644 index 000000000..11f5a64f8 --- /dev/null +++ b/db/migrations/20260618170000_game_stats_snapshot_partition.sql @@ -0,0 +1,39 @@ +-- Convert `game_stats_snapshot` to a monthly range-partitioned table. +-- +-- This migration only handles the structural change (rename old table, create the new +-- partitioned table, configure pg_partman). The data copy is a separate migration +-- (20260618180000) so that the two steps can be reasoned about independently and the +-- data migration can be re-run or batched manually on prod if the table is very large. +-- +-- pg_partman creates partitions centred on NOW() (+/- p_premake months). Historical +-- rows inserted by the copy migration land in the DEFAULT partition; run +-- SELECT partman.partition_data_time('public.game_stats_snapshot'); +-- on prod after deploying to move them into proper monthly partitions. This avoids +-- pre-creating ~55 empty historical partitions in every CI/fresh-install DB. + +ALTER TABLE game_stats_snapshot RENAME TO game_stats_snapshot_old; + +CREATE TABLE game_stats_snapshot ( + user_id BIGINT REFERENCES account(id) NOT NULL, + game GAME NOT NULL, + timestamp TIMESTAMPTZ NOT NULL, + playcount BIGINT NOT NULL, + ratings JSONB NOT NULL, + classes JSONB NOT NULL, + rankings JSONB NOT NULL, + PRIMARY KEY (user_id, game, timestamp) +) PARTITION BY RANGE (timestamp); + +CREATE INDEX ON game_stats_snapshot (game, timestamp DESC); + +SELECT partman.create_parent( + p_parent_table => 'public.game_stats_snapshot', + p_control => 'timestamp', + p_interval => '1 month', + p_premake => 3 +); + +-- All historical data is valuable — no automatic partition drops. +UPDATE partman.part_config +SET retention = NULL +WHERE parent_table = 'public.game_stats_snapshot'; diff --git a/db/migrations/20260618180000_game_stats_snapshot_copy.sql b/db/migrations/20260618180000_game_stats_snapshot_copy.sql new file mode 100644 index 000000000..9c207a3a5 --- /dev/null +++ b/db/migrations/20260618180000_game_stats_snapshot_copy.sql @@ -0,0 +1,24 @@ +-- Copy data from the unpartitioned game_stats_snapshot_old into the new partitioned +-- game_stats_snapshot table created by migration 20260618170000. +-- +-- On prod this covers ~4.7 M rows (1.3 GB). The INSERT is safe to run in the migration +-- runner; on very large prod instances it can also be done manually in batches with +-- partman.partition_data_time() before deploying, then just DROP the old table here. + +INSERT INTO game_stats_snapshot SELECT * FROM game_stats_snapshot_old; + +DROP TABLE game_stats_snapshot_old; + +-- Redistribute all rows from the DEFAULT partition into their correct monthly partitions. +-- partition_data_time() creates missing monthly partitions on the fly and returns the +-- number of rows moved; loop until it returns 0. +DO $$ +DECLARE + moved BIGINT; +BEGIN + LOOP + SELECT partman.partition_data_time('public.game_stats_snapshot', p_batch_count => 1000) + INTO moved; + EXIT WHEN moved = 0; + END LOOP; +END $$; diff --git a/dev/postgres-init.sql b/dev/postgres-init.sql index 3b3f22b9c..b294de56b 100644 --- a/dev/postgres-init.sql +++ b/dev/postgres-init.sql @@ -4,6 +4,8 @@ GRANT ALL PRIVILEGES ON DATABASE tachi_dev TO tachi; -- Match genesis: pg_stat_statements (requires shared_preload_libraries in docker-compose). \connect tachi_dev CREATE EXTENSION IF NOT EXISTS pg_stat_statements; +CREATE SCHEMA IF NOT EXISTS partman; +CREATE EXTENSION IF NOT EXISTS pg_partman SCHEMA partman; -- Read-only user for Grafana dashboards. DO $$ diff --git a/docker-compose-dev.yml b/docker-compose-dev.yml index 609a85d47..7b0782062 100644 --- a/docker-compose-dev.yml +++ b/docker-compose-dev.yml @@ -10,17 +10,18 @@ services: - tachi-redis:/data tachi-postgres: container_name: tachi-postgres - image: postgres:18 + image: ghcr.io/zkldi/zk-postgres:18@sha256:484fa9e6501ec93c5729ba9f9312d8ee4dfe0c136f1b4881fab5e0aed5aa68fd restart: unless-stopped ports: - "5432:5432" - # Dev observability: pg_stat_statements + slow-query logging + auto_explain plans. - # Recreate the container (or bump volume) after changing these; existing DBs: run - # CREATE EXTENSION pg_stat_statements; once if missing. command: - postgres - -c - - shared_preload_libraries=pg_stat_statements,auto_explain + - shared_preload_libraries=pg_stat_statements,auto_explain,pg_partman_bgw + - -c + - pg_partman_bgw.interval=3600 + - -c + - pg_partman_bgw.dbname=tachi - -c - pg_stat_statements.track=all - -c @@ -47,7 +48,7 @@ services: # we want for tests. tachi-postgres-test: container_name: tachi-postgres-test - image: postgres:18 + image: ghcr.io/zkldi/zk-postgres:18@sha256:484fa9e6501ec93c5729ba9f9312d8ee4dfe0c136f1b4881fab5e0aed5aa68fd restart: unless-stopped ports: - "5433:5432" @@ -89,7 +90,7 @@ services: - -c - autovacuum_analyze_scale_factor=0.1 - -c - - shared_preload_libraries=pg_stat_statements + - shared_preload_libraries=pg_stat_statements,pg_partman_bgw - -c - pg_stat_statements.track=all # Sized for ~16 concurrent worker DBs + template + WAL headroom. @@ -103,24 +104,6 @@ services: POSTGRES_DB: postgres PGDATA: /var/lib/postgresql/data/pgdata - # PostgreSQL admin UI (http://localhost:5050). Dev + test servers are pre-registered - # on first launch via dev/pgadmin/servers.json (wipe tachi-pgadmin volume to reload). - tachi-pgadmin: - container_name: tachi-pgadmin - image: dpage/pgadmin4 - restart: unless-stopped - depends_on: - - tachi-postgres - - tachi-postgres-test - ports: - - "5050:80" - environment: - PGADMIN_DEFAULT_EMAIL: pgadmin@example.com - PGADMIN_DEFAULT_PASSWORD: password - volumes: - - tachi-pgadmin:/var/lib/pgadmin - - "./dev/pgadmin/servers.json:/pgadmin4/servers.json:ro" - tachi-s3: container_name: tachi-s3 image: quay.io/minio/minio:RELEASE.2024-10-29T16-01-48Z @@ -134,7 +117,7 @@ services: - tachi-minio:/data entrypoint: /usr/bin/minio server /data --console-address=':9001' - # Local SMTP capture + web UI (http://localhost:8025). Server uses tachi-mailpit:1025 via env. + # Read emails on port 1025. tachi-mailpit: container_name: tachi-mailpit image: axllent/mailpit:v1.29.7 @@ -167,7 +150,7 @@ services: - /tachi/node_modules/ - /tachi/.bun - # Local observability (Grafana + Prometheus + Alloy). Grafana UI: http://localhost:3005 + # Local observability stack (Grafana + Prometheus + Alloy). Grafana UI: http://localhost:3005 # Alloy scrapes the API metrics endpoint at tachi-dev:9779 and remote-writes into Prometheus. tachi-grafana: container_name: tachi-grafana @@ -227,7 +210,6 @@ volumes: tachi-redis: tachi-logs: tachi-postgres: - tachi-pgadmin: tachi-minio: tachi-grafana: tachi-prometheus: \ No newline at end of file diff --git a/typescript/db/src/generated/index.ts b/typescript/db/src/generated/index.ts index 6cf221637..b83a45d6d 100644 --- a/typescript/db/src/generated/index.ts +++ b/typescript/db/src/generated/index.ts @@ -46,7 +46,6 @@ export { type default as AccountFollowingTable, type AccountFollowing, type NewA export { type import_id, type default as ImportTable, type Import, type NewImport, type ImportUpdate } from './public/Import'; export { type folder_id, type default as FolderTable, type Folder, type NewFolder, type FolderUpdate } from './public/Folder'; export { type default as PrivAccountCredentialTable, type PrivAccountCredential, type NewPrivAccountCredential, type PrivAccountCredentialUpdate } from './public/PrivAccountCredential'; -export { type score_id, type default as ScoreTable, type Score, type NewScore, type ScoreUpdate } from './public/Score'; export { type default as ImportTimingTable, type ImportTiming, type NewImportTiming, type ImportTimingUpdate } from './public/ImportTiming'; export { type default as QuestSubTable, type QuestSub, type NewQuestSub, type QuestSubUpdate } from './public/QuestSub'; export { type orphan_score_row_id, type default as OrphanScoreTable, type OrphanScore, type NewOrphanScore, type OrphanScoreUpdate } from './public/OrphanScore'; @@ -55,6 +54,7 @@ export { type table_id, type default as TableTable, type Table, type NewTable, t 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 raw_score_id, type default as RawScoreTable, type RawScore, type NewRawScore, type RawScoreUpdate } from './public/RawScore'; export { type priv_api_client_client_id, type default as PrivApiClientTable, type PrivApiClient, type NewPrivApiClient, type PrivApiClientUpdate } from './public/PrivApiClient'; export { type priv_api_token_token, type default as PrivApiTokenTable, type PrivApiToken, type NewPrivApiToken, type PrivApiTokenUpdate } from './public/PrivApiToken'; export { type default as SvcKshookSv6cSettingsTable, type SvcKshookSv6cSettings, type NewSvcKshookSv6cSettings, type SvcKshookSv6cSettingsUpdate } from './public/SvcKshookSv6cSettings'; @@ -70,6 +70,7 @@ 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 ScoreTable, type Score } from './public/Score'; export { type default as AuthLevel } from './public/AuthLevel'; export { type default as AccountBadgeKind } from './public/AccountBadgeKind'; export { type default as ImportStatus } from './public/ImportStatus'; diff --git a/typescript/db/src/generated/public/PbComposedFrom.ts b/typescript/db/src/generated/public/PbComposedFrom.ts index ce893eb67..e0f569505 100644 --- a/typescript/db/src/generated/public/PbComposedFrom.ts +++ b/typescript/db/src/generated/public/PbComposedFrom.ts @@ -2,14 +2,14 @@ // This file is automatically generated by Kanel. Do not modify manually. import type { pb_row_id } from './Pb'; -import type { score_id } from './Score'; +import type { raw_score_id } from './RawScore'; import type { ColumnType, Selectable, Insertable, Updateable } from 'kysely'; /** Represents the table public.pb_composed_from */ export default interface PbComposedFromTable { pb_id: ColumnType; - score_id: ColumnType; + score_id: ColumnType; merge_name: ColumnType; } diff --git a/typescript/db/src/generated/public/PublicSchema.ts b/typescript/db/src/generated/public/PublicSchema.ts index c8cbc25d7..7d910e8c7 100644 --- a/typescript/db/src/generated/public/PublicSchema.ts +++ b/typescript/db/src/generated/public/PublicSchema.ts @@ -46,7 +46,6 @@ import type { default as AccountFollowingTable } from './AccountFollowing'; import type { default as ImportTable } from './Import'; import type { default as FolderTable } from './Folder'; import type { default as PrivAccountCredentialTable } from './PrivAccountCredential'; -import type { default as ScoreTable } from './Score'; import type { default as ImportTimingTable } from './ImportTiming'; import type { default as QuestSubTable } from './QuestSub'; import type { default as OrphanScoreTable } from './OrphanScore'; @@ -55,6 +54,7 @@ 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 RawScoreTable } from './RawScore'; import type { default as PrivApiClientTable } from './PrivApiClient'; import type { default as PrivApiTokenTable } from './PrivApiToken'; import type { default as SvcKshookSv6cSettingsTable } from './SvcKshookSv6cSettings'; @@ -70,6 +70,7 @@ 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 ScoreTable } from './Score'; export default interface PublicSchema { priv_svc_cg_card_info: PrivSvcCgCardInfoTable; @@ -162,8 +163,6 @@ export default interface PublicSchema { priv_account_credential: PrivAccountCredentialTable; - score: ScoreTable; - import_timing: ImportTimingTable; quest_sub: QuestSubTable; @@ -180,6 +179,8 @@ export default interface PublicSchema { job_queue: JobQueueTable; + raw_score: RawScoreTable; + priv_api_client: PrivApiClientTable; priv_api_token: PrivApiTokenTable; @@ -209,4 +210,6 @@ export default interface PublicSchema { import_class: ImportClassTable; game_stats_snapshot: GameStatsSnapshotTable; + + score: ScoreTable; } diff --git a/typescript/db/src/generated/public/RawScore.ts b/typescript/db/src/generated/public/RawScore.ts new file mode 100644 index 000000000..3e63057e6 --- /dev/null +++ b/typescript/db/src/generated/public/RawScore.ts @@ -0,0 +1,51 @@ +// @generated +// This file is automatically generated by Kanel. Do not modify manually. + +import type { account_id } from './Account'; +import type { chart_id } from './Chart'; +import type { default as Game } from './Game'; +import type { ColumnType, Selectable, Insertable, Updateable } from 'kysely'; + +/** Identifier type for public.raw_score */ +export type raw_score_id = string; + +/** Represents the table public.raw_score */ +export default interface RawScoreTable { + id: ColumnType; + + user_id: ColumnType; + + chart_id: ColumnType; + + game: ColumnType; + + session_id: ColumnType; + + import_id: ColumnType; + + data: ColumnType; + + derived_data: ColumnType; + + calculated_data: ColumnType; + + judgements: ColumnType; + + meta: ColumnType; + + time_achieved: ColumnType; + + time_added: ColumnType; + + highlight: ColumnType; + + comment: ColumnType; + + committed: ColumnType; +} + +export type RawScore = Selectable; + +export type NewRawScore = Insertable; + +export type RawScoreUpdate = Updateable; diff --git a/typescript/db/src/generated/public/Score.ts b/typescript/db/src/generated/public/Score.ts index fad893c8c..a6050716b 100644 --- a/typescript/db/src/generated/public/Score.ts +++ b/typescript/db/src/generated/public/Score.ts @@ -1,51 +1,45 @@ // @generated // This file is automatically generated by Kanel. Do not modify manually. +import type { raw_score_id } from './RawScore'; import type { account_id } from './Account'; import type { chart_id } from './Chart'; import type { default as Game } from './Game'; -import type { ColumnType, Selectable, Insertable, Updateable } from 'kysely'; +import type { ColumnType, Selectable } from 'kysely'; -/** Identifier type for public.score */ -export type score_id = string; - -/** Represents the table public.score */ +/** Represents the view public.score */ export default interface ScoreTable { - id: ColumnType; + id: ColumnType; - user_id: ColumnType; + user_id: ColumnType; - chart_id: ColumnType; + chart_id: ColumnType; - game: ColumnType; + game: ColumnType; - session_id: ColumnType; + session_id: ColumnType; - import_id: ColumnType; + import_id: ColumnType; - data: ColumnType; + data: ColumnType; - derived_data: ColumnType; + derived_data: ColumnType; - calculated_data: ColumnType; + calculated_data: ColumnType; - judgements: ColumnType; + judgements: ColumnType; - meta: ColumnType; + meta: ColumnType; - time_achieved: ColumnType; + time_achieved: ColumnType; - time_added: ColumnType; + time_added: ColumnType; - highlight: ColumnType; + highlight: ColumnType; - comment: ColumnType; + comment: ColumnType; - committed: ColumnType; + committed: ColumnType; } export type Score = Selectable; - -export type NewScore = Insertable; - -export type ScoreUpdate = Updateable; diff --git a/typescript/server/src/game-implementations/utils/pb-merge.ts b/typescript/server/src/game-implementations/utils/pb-merge.ts index a3f563180..479ab3510 100644 --- a/typescript/server/src/game-implementations/utils/pb-merge.ts +++ b/typescript/server/src/game-implementations/utils/pb-merge.ts @@ -61,7 +61,7 @@ export function CreatePBMergeFor( applicator: (base: PBScoreDocumentNoRank, score: ScoreDocument) => void, ): PBMergeFunction { return async (userID, chartID, asOfTimestamp, base) => { - let q = DB.selectFrom("score") + let q = DB.selectFrom("raw_score as score") .innerJoin("chart", "chart.id", "score.chart_id") .innerJoin("song", "song.id", "chart.song_id") .leftJoin("import", "import.id", "score.import_id") diff --git a/typescript/server/src/lib/db-formats/score.ts b/typescript/server/src/lib/db-formats/score.ts index 62913eee5..cd6fdb5c4 100644 --- a/typescript/server/src/lib/db-formats/score.ts +++ b/typescript/server/src/lib/db-formats/score.ts @@ -107,7 +107,7 @@ export async function LoadScoreDocumentById(scoreID: string): Promise> { - const rows = await DB.selectFrom("score") + const rows = await DB.selectFrom("raw_score as score") .innerJoin("chart", "chart.id", "score.chart_id") .innerJoin("song", "song.id", "chart.song_id") .leftJoin("import", "import.id", "score.import_id") diff --git a/typescript/server/src/lib/dirty-queues/dirty-queue-race.test.ts b/typescript/server/src/lib/dirty-queues/dirty-queue-race.test.ts index 776c2b314..50d696957 100644 --- a/typescript/server/src/lib/dirty-queues/dirty-queue-race.test.ts +++ b/typescript/server/src/lib/dirty-queues/dirty-queue-race.test.ts @@ -110,7 +110,7 @@ async function insertCommittedScore(opts: { const { data, derived, judgements } = mongoScoreDataToPg("iidx-sp", doc.scoreData); const ts = UnixMillisecondsToISO8601(1_700_000_000_000); - await DB.insertInto("score") + await DB.insertInto("raw_score") .values({ id: opts.scoreId, user_id: opts.userId, @@ -155,7 +155,7 @@ async function insertUncommittedScore(opts: { const { data, derived, judgements } = mongoScoreDataToPg("iidx-sp", doc.scoreData); const ts = UnixMillisecondsToISO8601(1_700_000_100_000); - await DB.insertInto("score") + await DB.insertInto("raw_score") .values({ id: opts.scoreId, user_id: opts.userId, @@ -478,9 +478,9 @@ describe("dirty queue race harness", () => { await drainPbDirty(); expect(await readPbPercent(userId, chartId)).toBe(50); - await DB.updateTable("score") + await DB.updateTable("raw_score") .set({ committed: true }) - .where("score.import_id", "=", importId) + .where("raw_score.import_id", "=", importId) .execute(); const afterCommitDirty = await DB.selectFrom("pb_dirty") diff --git a/typescript/server/src/lib/imports/imports.test.ts b/typescript/server/src/lib/imports/imports.test.ts index ec20fb820..17c523989 100644 --- a/typescript/server/src/lib/imports/imports.test.ts +++ b/typescript/server/src/lib/imports/imports.test.ts @@ -66,7 +66,7 @@ async function insertIidxScore(opts: { const { data, derived, judgements } = mongoScoreDataToPg("iidx-sp", doc.scoreData); const ts = UnixMillisecondsToISO8601(timeMs); - await DB.insertInto("score") + await DB.insertInto("raw_score") .values({ id: opts.scoreId, user_id: opts.userId, @@ -173,15 +173,15 @@ describe("RevertImport", () => { .executeTakeFirst(); expect(importRow).toBeUndefined(); - const s1 = await DB.selectFrom("score") + const s1 = await DB.selectFrom("raw_score") .select("id") .where("id", "=", "score_1") .executeTakeFirst(); - const s2 = await DB.selectFrom("score") + const s2 = await DB.selectFrom("raw_score") .select("id") .where("id", "=", "score_2") .executeTakeFirst(); - const s3 = await DB.selectFrom("score") + const s3 = await DB.selectFrom("raw_score") .select("id") .where("id", "=", "score_3") .executeTakeFirst(); @@ -265,7 +265,7 @@ describe("RevertImport", () => { .executeTakeFirst(); expect(link).toBeUndefined(); - const deletedScore = await DB.selectFrom("score") + const deletedScore = await DB.selectFrom("raw_score") .select("id") .where("id", "=", scoreId) .executeTakeFirst(); diff --git a/typescript/server/src/lib/jobs/drain-dirty-queues.ts b/typescript/server/src/lib/jobs/drain-dirty-queues.ts index 88a8fd104..171c9bbcd 100644 --- a/typescript/server/src/lib/jobs/drain-dirty-queues.ts +++ b/typescript/server/src/lib/jobs/drain-dirty-queues.ts @@ -149,7 +149,7 @@ export async function drainSessionDirty(): Promise { const runStartedAt = await newCalculationRunStartedAt(); // eslint-disable-next-line no-await-in-loop - const scoreRows = await DB.selectFrom("score") + const scoreRows = await DB.selectFrom("raw_score as score") .innerJoin("chart", "chart.id", "score.chart_id") .innerJoin("song", "song.id", "chart.song_id") .leftJoin("import", "import.id", "score.import_id") diff --git a/typescript/server/src/lib/score-import/framework/pb/create-pb-doc.ts b/typescript/server/src/lib/score-import/framework/pb/create-pb-doc.ts index c854250b6..de816af64 100644 --- a/typescript/server/src/lib/score-import/framework/pb/create-pb-doc.ts +++ b/typescript/server/src/lib/score-import/framework/pb/create-pb-doc.ts @@ -48,7 +48,7 @@ async function GetBaseScoreForPB( ): Promise { const sortVal = defaultMetricSortValueSql(game); - let q = DB.selectFrom("score") + let q = DB.selectFrom("raw_score as score") .innerJoin("chart", "chart.id", "score.chart_id") .innerJoin("song", "song.id", "chart.song_id") .leftJoin("import", "import.id", "score.import_id") diff --git a/typescript/server/src/lib/score-import/framework/pb/pb-dirty.test.ts b/typescript/server/src/lib/score-import/framework/pb/pb-dirty.test.ts index b164033e2..8a9cc8877 100644 --- a/typescript/server/src/lib/score-import/framework/pb/pb-dirty.test.ts +++ b/typescript/server/src/lib/score-import/framework/pb/pb-dirty.test.ts @@ -70,7 +70,7 @@ async function insertScoreRow(opts: { const { data, derived, judgements } = mongoScoreDataToPg("iidx-sp", doc.scoreData); const ts = UnixMillisecondsToISO8601(1_700_000_000_000); - await DB.insertInto("score") + await DB.insertInto("raw_score") .values({ id: opts.scoreId, user_id: opts.userId, @@ -147,9 +147,9 @@ describe("pb_dirty trigger and last_clean_started_at guards", () => { importId, }); - await DB.updateTable("score") + await DB.updateTable("raw_score") .set({ committed: true }) - .where("score.id", "=", "score_staged") + .where("raw_score.id", "=", "score_staged") .execute(); const dirty = await DB.selectFrom("pb_dirty") diff --git a/typescript/server/src/lib/score-import/framework/pg/ensure-import-stub.ts b/typescript/server/src/lib/score-import/framework/pg/ensure-import-stub.ts index f02a8004e..983791957 100644 --- a/typescript/server/src/lib/score-import/framework/pg/ensure-import-stub.ts +++ b/typescript/server/src/lib/score-import/framework/pg/ensure-import-stub.ts @@ -32,14 +32,47 @@ export async function ensureImportStub( } /** - * Removes an in-progress import run: staged scores and the import stub (and dependent rows cascade or explicit deletes). + * Removes an in-progress import run: staged scores and the import stub (dependent rows cascade). */ export async function deleteImportRun(importId: string): Promise { - await DB.deleteFrom("score") - .where("import_id", "=", importId) - .where("committed", "=", false) + await DB.deleteFrom("raw_score as score") + .where("score.import_id", "=", importId) + .where("score.committed", "=", false) .execute(); // import_* / import_timing / import_game rows cascade. orphan_score.import_id is set null (not deleted). await DB.deleteFrom("import").where("id", "=", importId).execute(); } + +/** + * Cleans up all uncommitted scores left over from crashed/stuck imports for a user. + * Called at the start of every new import (after acquiring the lock) so that orphaned + * committed=false rows from a previous crash never block re-importing or pollute PBs. + * + * Deletes via the `import` table first (in_progress stubs → their staged scores), then + * sweeps directly for any remaining committed=false rows that have no matching import row. + */ +export async function cleanUpStaleImportsForUser( + userId: integer, + currentImportId: string, +): Promise { + const staleImports = await DB.selectFrom("import") + .select("import.id") + .where("import.user_id", "=", userId) + .where("import.status", "=", "in_progress") + .execute(); + + for (const stale of staleImports) { + if (stale.id !== currentImportId) { + // eslint-disable-next-line no-await-in-loop + await deleteImportRun(stale.id); + } + } + + // Safety net: delete any committed=false scores whose import row was already removed. + await DB.deleteFrom("raw_score as score") + .where("score.user_id", "=", userId) + .where("score.committed", "=", false) + .where("score.import_id", "!=", currentImportId) + .execute(); +} diff --git a/typescript/server/src/lib/score-import/framework/pg/finalize-import-pg.ts b/typescript/server/src/lib/score-import/framework/pg/finalize-import-pg.ts index 334ef3ef0..5dd6e17bf 100644 --- a/typescript/server/src/lib/score-import/framework/pg/finalize-import-pg.ts +++ b/typescript/server/src/lib/score-import/framework/pg/finalize-import-pg.ts @@ -88,10 +88,10 @@ export async function finalizeImportToPostgres( .execute(); await db - .updateTable("score") + .updateTable("raw_score as score") .set({ committed: true }) - .where("import_id", "=", importID) - .where("committed", "=", false) + .where("score.import_id", "=", importID) + .where("score.committed", "=", false) .execute(); await db.deleteFrom("import_game").where("id", "=", importID).execute(); diff --git a/typescript/server/src/lib/score-import/framework/pg/mongo-score-to-pg.ts b/typescript/server/src/lib/score-import/framework/pg/mongo-score-to-pg.ts index 56e4f9b32..76b99d10d 100644 --- a/typescript/server/src/lib/score-import/framework/pg/mongo-score-to-pg.ts +++ b/typescript/server/src/lib/score-import/framework/pg/mongo-score-to-pg.ts @@ -1,4 +1,4 @@ -import type { NewScore } from "tachi-db"; +import type { NewRawScore } from "tachi-db"; import { mongoScoreDataToPg } from "#lib/v3/migration-tools"; import { UnixMillisecondsToISO8601 } from "#utils/time"; @@ -16,7 +16,7 @@ export function mongoScoreDocumentToNewScoreRow( importId: string | null; sessionId: string | null; }, -): NewScore { +): NewRawScore { const game = score.game; const { data, derived, judgements } = mongoScoreDataToPg(game, score.scoreData); diff --git a/typescript/server/src/lib/score-import/framework/pg/score-visibility.ts b/typescript/server/src/lib/score-import/framework/pg/score-visibility.ts index 6ad1719d5..8fc1ccd7f 100644 --- a/typescript/server/src/lib/score-import/framework/pg/score-visibility.ts +++ b/typescript/server/src/lib/score-import/framework/pg/score-visibility.ts @@ -6,7 +6,8 @@ import { type ExpressionBuilder, sql } from "kysely"; import { getActiveImportId } from "../import-run-context"; /** - * SQL predicate: score row is visible for import logic (committed scores, or scores from the active import run). + * SQL predicate for queries against `raw_score as score`: visible means committed, + * or belonging to the currently active import run (staged scores). */ export function scoreVisiblePredicate(eb: ExpressionBuilder) { const importId = getActiveImportId(); @@ -19,10 +20,8 @@ export function scoreVisiblePredicate(eb: ExpressionBuilder) } /** - * Same as {@link scoreVisiblePredicate}, but as a standalone expression for queries that join `score` with other tables (Kysely's `ExpressionBuilder` table set differs). - * - * TODO(zk): move scores into an "internal_scores" table, then have a "scores" view which contains only - * the visible ones - allows us to ban users and stuff. + * Same as {@link scoreVisiblePredicate}, but as a standalone SQL fragment for queries + * that join `raw_score as score` with other tables (Kysely's ExpressionBuilder table set differs). */ export function scoreVisibleSql() { const importId = getActiveImportId(); @@ -36,17 +35,17 @@ export function scoreVisibleSql() { /** Deletes all uncommitted scores for an import run (failed import cleanup). */ export async function deleteUncommittedScoresForImport(importId: string): Promise { - await DB.deleteFrom("score") - .where("import_id", "=", importId) - .where("committed", "=", false) + await DB.deleteFrom("raw_score as score") + .where("score.import_id", "=", importId) + .where("score.committed", "=", false) .execute(); } /** Marks staged scores as committed after successful post-import steps. */ export async function commitScoresForImport(importId: string): Promise { - await DB.updateTable("score") + await DB.updateTable("raw_score as score") .set({ committed: true }) - .where("import_id", "=", importId) - .where("committed", "=", false) + .where("score.import_id", "=", importId) + .where("score.committed", "=", false) .execute(); } diff --git a/typescript/server/src/lib/score-import/framework/score-importing/insert-score.ts b/typescript/server/src/lib/score-import/framework/score-importing/insert-score.ts index 618a14d01..ad5811af5 100644 --- a/typescript/server/src/lib/score-import/framework/score-importing/insert-score.ts +++ b/typescript/server/src/lib/score-import/framework/score-importing/insert-score.ts @@ -67,7 +67,7 @@ export async function InsertQueue(userID: integer) { }), ); - await DB.insertInto("score").values(rows).execute(); + await DB.insertInto("raw_score").values(rows).execute(); } catch (err) { log.warn( { err }, diff --git a/typescript/server/src/lib/score-import/framework/score-importing/score-import-main.ts b/typescript/server/src/lib/score-import/framework/score-importing/score-import-main.ts index 39d2a9e95..2c81c457b 100644 --- a/typescript/server/src/lib/score-import/framework/score-importing/score-import-main.ts +++ b/typescript/server/src/lib/score-import/framework/score-importing/score-import-main.ts @@ -9,6 +9,7 @@ import { LoadImportDocumentById } from "#lib/db-formats/import-document"; import { clearPbDirtyForUser } from "#lib/jobs/drain-dirty-queues"; import { runWithImportContext } from "#lib/score-import/framework/import-run-context"; import { + cleanUpStaleImportsForUser, deleteImportRun, ensureImportStub, } from "#lib/score-import/framework/pg/ensure-import-stub"; @@ -97,6 +98,8 @@ export default async function ScoreImportMain( } return runWithImportContext(importID, async () => { + // Wipe any committed=false rows from previous crashed imports before we begin. + await cleanUpStaleImportsForUser(user.id, importID); await deleteImportRun(importID); const timeStarted = Date.now(); diff --git a/typescript/server/src/lib/score-import/framework/score-importing/score-importing.test.ts b/typescript/server/src/lib/score-import/framework/score-importing/score-importing.test.ts index 6be8b653c..05b67c5f2 100644 --- a/typescript/server/src/lib/score-import/framework/score-importing/score-importing.test.ts +++ b/typescript/server/src/lib/score-import/framework/score-importing/score-importing.test.ts @@ -107,7 +107,7 @@ describe("score-importing duplicate score ID regression", () => { expect(flushAfterSkip).not.toBeNull(); expect(flushAfterSkip).toBe(0); - const rowCount = await DB.selectFrom("score") + const rowCount = await DB.selectFrom("raw_score") .select((eb) => eb.fn.countAll().as("c")) .executeTakeFirst() .then((r) => Number(r?.c ?? 0)); @@ -131,7 +131,7 @@ describe("score-importing duplicate score ID regression", () => { expect(results).toHaveLength(1); expect(results[0]?.success).toBe(true); - const rowCount = await DB.selectFrom("score") + const rowCount = await DB.selectFrom("raw_score") .select((eb) => eb.fn.countAll().as("c")) .executeTakeFirst() .then((r) => Number(r?.c ?? 0)); diff --git a/typescript/server/src/lib/score-import/framework/score-importing/score-importing.ts b/typescript/server/src/lib/score-import/framework/score-importing/score-importing.ts index 1b213740d..9c3dad25a 100644 --- a/typescript/server/src/lib/score-import/framework/score-importing/score-importing.ts +++ b/typescript/server/src/lib/score-import/framework/score-importing/score-importing.ts @@ -366,7 +366,7 @@ async function HydrateCheckAndInsertScore( } // committed or not - don't import scores twice. - const existing = await DB.selectFrom("score") + const existing = await DB.selectFrom("raw_score") .select("id") .where("id", "=", scoreID) .executeTakeFirst(); diff --git a/typescript/server/src/lib/score-import/framework/sessions/sessions.ts b/typescript/server/src/lib/score-import/framework/sessions/sessions.ts index f193e5b70..2a4f5a66d 100644 --- a/typescript/server/src/lib/score-import/framework/sessions/sessions.ts +++ b/typescript/server/src/lib/score-import/framework/sessions/sessions.ts @@ -300,9 +300,9 @@ export async function LoadScoresIntoSessions( .where("id", "=", session.sessionID) .execute(); - await DB.updateTable("score") + await DB.updateTable("raw_score as score") .set({ session_id: session.sessionID }) - .where("id", "in", scoreIDs) + .where("score.id", "in", scoreIDs) .execute(); } else { log.debug( @@ -330,9 +330,9 @@ export async function LoadScoresIntoSessions( }) .execute(); - await DB.updateTable("score") + await DB.updateTable("raw_score as score") .set({ session_id: session.sessionID }) - .where("id", "in", scoreIDs) + .where("score.id", "in", scoreIDs) .execute(); } diff --git a/typescript/server/src/lib/score-mutation/delete-scores.ts b/typescript/server/src/lib/score-mutation/delete-scores.ts index 5a1f3818d..2671e21ab 100644 --- a/typescript/server/src/lib/score-mutation/delete-scores.ts +++ b/typescript/server/src/lib/score-mutation/delete-scores.ts @@ -34,7 +34,7 @@ export async function DeleteMultipleScores( const scoreIDs = scores.map((e) => e.scoreID); const chartIDs = scores.map((e) => e.chartID); - const sessionLinks = await DB.selectFrom("score") + const sessionLinks = await DB.selectFrom("raw_score as score") .select("session_id") .where("id", "in", scoreIDs) .execute(); @@ -49,10 +49,10 @@ export async function DeleteMultipleScores( await DB.deleteFrom("pb_composed_from").where("score_id", "in", scoreIDs).execute(); - await DB.deleteFrom("score").where("id", "in", scoreIDs).execute(); + await DB.deleteFrom("raw_score").where("id", "in", scoreIDs).execute(); for (const sessionId of sessionIds) { - const remainingRows = await DB.selectFrom("score") + const remainingRows = await DB.selectFrom("raw_score as score") .innerJoin("chart", "chart.id", "score.chart_id") .innerJoin("song", "song.id", "chart.song_id") .leftJoin("import", "import.id", "score.import_id")