From 82076ea7675cdc73304858aefe5f5bd3c8c73d1a Mon Sep 17 00:00:00 2001 From: Daniel Salazar Date: Tue, 1 Sep 2026 22:04:32 -0700 Subject: [PATCH] feat: notification retention sweep (PUT-1668) (#3672) --- config.template.jsonc | 5 + .../database/SqliteDatabaseClient.test.ts | 9 +- .../clients/database/SqliteDatabaseClient.ts | 1 + .../migrations/mysql/mysql_mig_27.sql | 42 +++++++ .../migrations/postgres/postgres_mig_16.sql | 22 ++++ .../sqlite/0073_notification-created-at.sql | 23 ++++ .../notification/NotificationService.test.ts | 114 +++++++++++++++++- .../notification/NotificationService.ts | 58 +++++++++ .../stores/notification/NotificationStore.js | 35 ++++++ .../notification/NotificationStore.test.js | 84 +++++++++++++ src/backend/types.ts | 7 ++ 11 files changed, 398 insertions(+), 2 deletions(-) create mode 100644 src/backend/clients/database/migrations/mysql/mysql_mig_27.sql create mode 100644 src/backend/clients/database/migrations/postgres/postgres_mig_16.sql create mode 100644 src/backend/clients/database/migrations/sqlite/0073_notification-created-at.sql diff --git a/config.template.jsonc b/config.template.jsonc index 67f042708..4d1fbf87e 100644 --- a/config.template.jsonc +++ b/config.template.jsonc @@ -273,6 +273,11 @@ // "emailBatchSeconds": 90 // }, + // ── Notifications ─────────────────────────────────────────────────── + // How long a notification is kept, in days from creation. Acknowledged or + // not, a row past this is swept. Set 0 to keep everything forever. + "notificationRetentionDays": 14, + // ── Alarms / alerting ─────────────────────────────────────────────── // Where system alarms go. Severity is the routing decision — each // transport takes everything at or above its own `minSeverity`: diff --git a/src/backend/clients/database/SqliteDatabaseClient.test.ts b/src/backend/clients/database/SqliteDatabaseClient.test.ts index 076633f63..d890246ec 100644 --- a/src/backend/clients/database/SqliteDatabaseClient.test.ts +++ b/src/backend/clients/database/SqliteDatabaseClient.test.ts @@ -27,7 +27,7 @@ import { DatabaseClientFactory } from './index.js'; import { SqliteDatabaseClient } from './SqliteDatabaseClient.js'; /** Highest schema version the migration table can reach. */ -const CURRENT_SCHEMA_VERSION = 68; +const CURRENT_SCHEMA_VERSION = 69; /** * These suites migrate real files on disk. Idle they finish in well under a @@ -367,6 +367,13 @@ describe('SqliteDatabaseClient — boot and migrations', { timeout: DISK_MIGRATI }); }); + it('indexes notification.created_at for the retention sweep', async () => { + const rows = await client.read('PRAGMA index_list(`notification`)'); + expect(rows.map((r) => String(r.name))).toContain( + 'idx_notification_created_at', + ); + }); + it('leaves an already-migrated database untouched on a second boot', async () => { const dir = mkdtempSync(join(tmpdir(), 'puter-sqlite-')); const path = join(dir, 'nested', 'puter.sqlite'); diff --git a/src/backend/clients/database/SqliteDatabaseClient.ts b/src/backend/clients/database/SqliteDatabaseClient.ts index 4267e6386..4a37f0737 100644 --- a/src/backend/clients/database/SqliteDatabaseClient.ts +++ b/src/backend/clients/database/SqliteDatabaseClient.ts @@ -102,6 +102,7 @@ const AVAILABLE_MIGRATIONS: [number, string[]][] = [ [65, ['0070_drop-orphaned-default-groups.sql']], [66, ['0071_share_issuer_index.sql']], [67, ['0072_notification-scope.sql']], + [68, ['0073_notification-created-at.sql']], ]; export class SqliteDatabaseClient extends AbstractDatabaseClient { diff --git a/src/backend/clients/database/migrations/mysql/mysql_mig_27.sql b/src/backend/clients/database/migrations/mysql/mysql_mig_27.sql new file mode 100644 index 000000000..e1f432bd7 --- /dev/null +++ b/src/backend/clients/database/migrations/mysql/mysql_mig_27.sql @@ -0,0 +1,42 @@ +-- Copyright (C) 2024-present Puter Technologies Inc. +-- +-- This file is part of Puter. +-- +-- Puter is free software: you can redistribute it and/or modify +-- it under the terms of the GNU Affero General Public License as published +-- by the Free Software Foundation, either version 3 of the License, or +-- (at your option) any later version. +-- +-- This program is distributed in the hope that it will be useful, +-- but WITHOUT ANY WARRANTY; without even the implied warranty of +-- MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +-- GNU Affero General Public License for more details. +-- +-- You should have received a copy of the GNU Affero General Public License +-- along with this program. If not, see . + +-- Index `notification`.`created_at`. The retention sweep selects by age alone +-- and by nothing else, so no existing index on the table narrows it — every +-- pass would otherwise scan the whole table to find the oldest few hundred +-- rows. +-- +-- Idempotent: the guarded procedure, as mig_26. There is no per-file +-- applied-state tracking, so a replay has to be a no-op. + +DROP PROCEDURE IF EXISTS _puter_add_notification_created_at_index; +DELIMITER // +CREATE PROCEDURE _puter_add_notification_created_at_index() +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM INFORMATION_SCHEMA.STATISTICS + WHERE TABLE_SCHEMA = DATABASE() + AND TABLE_NAME = 'notification' + AND INDEX_NAME = 'idx_notification_created_at' + ) THEN + ALTER TABLE `notification` ADD INDEX `idx_notification_created_at` + (`created_at`); + END IF; +END // +DELIMITER ; +CALL _puter_add_notification_created_at_index(); +DROP PROCEDURE IF EXISTS _puter_add_notification_created_at_index; diff --git a/src/backend/clients/database/migrations/postgres/postgres_mig_16.sql b/src/backend/clients/database/migrations/postgres/postgres_mig_16.sql new file mode 100644 index 000000000..68a9812cd --- /dev/null +++ b/src/backend/clients/database/migrations/postgres/postgres_mig_16.sql @@ -0,0 +1,22 @@ +-- Copyright (C) 2024-present Puter Technologies Inc. +-- +-- This file is part of Puter. +-- +-- Puter is free software: you can redistribute it and/or modify +-- it under the terms of the GNU Affero General Public License as published +-- by the Free Software Foundation, either version 3 of the License, or +-- (at your option) any later version. +-- +-- This program is distributed in the hope that it will be useful, +-- but WITHOUT ANY WARRANTY; without even the implied warranty of +-- MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +-- GNU Affero General Public License for more details. +-- +-- You should have received a copy of the GNU Affero General Public License +-- along with this program. If not, see . + +-- Index `notification`.`created_at` for the retention sweep. See +-- mysql/mysql_mig_27.sql for the rationale. + +CREATE INDEX IF NOT EXISTS idx_notification_created_at + ON notification (created_at); diff --git a/src/backend/clients/database/migrations/sqlite/0073_notification-created-at.sql b/src/backend/clients/database/migrations/sqlite/0073_notification-created-at.sql new file mode 100644 index 000000000..d5b0d002d --- /dev/null +++ b/src/backend/clients/database/migrations/sqlite/0073_notification-created-at.sql @@ -0,0 +1,23 @@ +-- Copyright (C) 2024-present Puter Technologies Inc. +-- +-- This file is part of Puter. +-- +-- Puter is free software: you can redistribute it and/or modify +-- it under the terms of the GNU Affero General Public License as published +-- by the Free Software Foundation, either version 3 of the License, or +-- (at your option) any later version. +-- +-- This program is distributed in the hope that it will be useful, +-- but WITHOUT ANY WARRANTY; without even the implied warranty of +-- MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +-- GNU Affero General Public License for more details. +-- +-- You should have received a copy of the GNU Affero General Public License +-- along with this program. If not, see . + +-- Index `notification`.`created_at` for the retention sweep. See +-- mysql/mysql_mig_27.sql for the rationale. A plain index add, so no table +-- rebuild. + +CREATE INDEX IF NOT EXISTS `idx_notification_created_at` + ON `notification` (`created_at`); diff --git a/src/backend/services/notification/NotificationService.test.ts b/src/backend/services/notification/NotificationService.test.ts index ca20daff6..749b86539 100644 --- a/src/backend/services/notification/NotificationService.test.ts +++ b/src/backend/services/notification/NotificationService.test.ts @@ -63,7 +63,11 @@ const makeUser = async (): Promise<{ id: number; username: string }> => { }; beforeAll(async () => { - server = await setupTestServer(); + // Retention is opt-in (config.default.json ships without it), so the + // sweepExpired suite below needs it turned on explicitly. + server = await setupTestServer({ + notificationRetentionDays: 14, + } as never); notifications = server.services .notification as unknown as NotificationService; }); @@ -486,3 +490,111 @@ describe('NotificationService — delivery receipts', () => { expect(rows).toHaveLength(0); }); }); + +describe('NotificationService.sweepExpired', () => { + /** Rows aged past the window, written straight in so the age is fixture. */ + const seedExpired = async ( + userId: number, + count: number, + days = 20, + ): Promise => { + const when = new Date(Date.now() - days * 86_400_000) + .toISOString() + .replace('T', ' ') + .slice(0, 19); + const uids = Array.from({ length: count }, () => uuidv4()); + await server.clients.db.batchWrite( + uids.map((uid) => ({ + statement: + 'INSERT INTO `notification` (`uid`, `user_id`, `value`, `created_at`) ' + + 'VALUES (?, ?, ?, ?)', + values: [uid, userId, '{}', when], + })), + ); + return uids; + }; + + it('removes what has aged out and leaves the mailbox otherwise intact', async () => { + await notifications.sweepExpired(); + const user = await makeUser(); + const [expired] = await seedExpired(user.id, 1); + const kept = await server.stores.notification.create({ + userId: user.id, + value: { title: 'still fresh' }, + }); + + expect(await notifications.sweepExpired()).toBe(1); + + expect(await server.stores.notification.getByUid(expired)).toBeNull(); + const unread = await server.stores.notification.listByUserId(user.id, { + onlyUnacknowledged: true, + }); + expect(unread.map((r) => r.uid)).toEqual([kept.uid]); + // Replay only carries what was never shown, and that is unchanged too. + const unseen = await server.stores.notification.listByUserId(user.id, { + filter: 'unseen', + }); + expect(unseen.map((r) => r.uid)).toEqual([kept.uid]); + }); + + it('keeps batching until the window is clean', async () => { + await notifications.sweepExpired(); + const user = await makeUser(); + // More than one batch takes, so the loop has to come back around. + await seedExpired(user.id, 600); + + expect(await notifications.sweepExpired()).toBe(600); + expect(await server.stores.notification.listByUserId(user.id)).toEqual( + [], + ); + // Nothing left, so the next pass ends on its first batch. + expect(await notifications.sweepExpired()).toBe(0); + }); + + it('stops at the pass cap and leaves the rest for the next sweep', async () => { + await notifications.sweepExpired(); + const user = await makeUser(); + // 50 passes * 500/batch = 25,000 — a backlog past that so the cap, + // not a short batch, is what ends the first call. + await seedExpired(user.id, 25_050); + + expect(await notifications.sweepExpired()).toBe(25_000); + expect( + await server.stores.notification.listByUserId(user.id), + ).toHaveLength(50); + expect(await notifications.sweepExpired()).toBe(50); + }); + + it('sweeps nothing when no retention is configured', async () => { + const unbounded = await setupTestServer({ + notificationRetentionDays: 0, + } as never); + try { + const service = unbounded.services + .notification as unknown as NotificationService; + const created = await unbounded.stores.user.create({ + username: `notif-${Math.random().toString(36).slice(2, 10)}`, + uuid: uuidv4(), + password: null, + email: `retention-${Date.now()}@test.local`, + requires_email_confirmation: false, + }); + const when = new Date(Date.now() - 400 * 86_400_000) + .toISOString() + .replace('T', ' ') + .slice(0, 19); + await unbounded.clients.db.write( + 'INSERT INTO `notification` (`uid`, `user_id`, `value`, `created_at`) ' + + 'VALUES (?, ?, ?, ?)', + [uuidv4(), created.id, '{}', when], + ); + + expect(await service.sweepExpired()).toBe(0); + expect( + await unbounded.stores.notification.listByUserId(created.id), + ).toHaveLength(1); + } finally { + await unbounded.shutdown(); + } + }); +}); diff --git a/src/backend/services/notification/NotificationService.ts b/src/backend/services/notification/NotificationService.ts index 4f5e3f867..5a568adce 100644 --- a/src/backend/services/notification/NotificationService.ts +++ b/src/backend/services/notification/NotificationService.ts @@ -36,6 +36,13 @@ export type { } from './notificationTypes.js'; export { canViewNotification } from './notificationAudience.js'; +/** How often the retention sweep runs. */ +const RETENTION_SWEEP_INTERVAL_MS = 60 * 60 * 1000; +/** Rows one delete takes. Small enough not to hold a lock anyone waits on. */ +const RETENTION_BATCH_SIZE = 500; +/** Batches one sweep takes, so a large backlog drains over several passes. */ +const RETENTION_MAX_BATCHES = 50; + /** * Notification orchestration — glues the NotificationStore (DB) to the event * bus (socket push) and handles lifecycle events (user connects → send unreads, @@ -49,8 +56,11 @@ export class NotificationService extends PuterService { #pendingWrites = new Map>(); /** User.id → debounce timeout */ #connectTimeouts = new Map>(); + #retentionSweep: ReturnType | null = null; override onServerStart(): void { + this.#armRetentionSweep(); + // When a user opens the GUI, send their pending unreads. this.clients.event.on( 'web.socket.user-connected', @@ -93,6 +103,11 @@ export class NotificationService extends PuterService { ); } + override onServerPrepareShutdown(): void { + if (this.#retentionSweep) clearInterval(this.#retentionSweep); + this.#retentionSweep = null; + } + // -- Public API -------------------------------------------------- /** @@ -251,6 +266,33 @@ export class NotificationService extends PuterService { ); } + /** + * Drop notifications past the retention window, in batches, and report how + * many went. Deleting is all there is to it: nothing is pushed, because a + * two-week-old row is not news, and a client listing again simply stops + * seeing it. + * + * Every node sweeps. Batches are small and the delete is idempotent, so two + * nodes overlapping costs a few empty batches, not correctness. + */ + async sweepExpired(): Promise { + const days = this.#retentionDays(); + if (days <= 0) return 0; + + let removed = 0; + for (let pass = 0; pass < RETENTION_MAX_BATCHES; pass++) { + const batch = await this.stores.notification.deleteCreatedBefore( + days, + RETENTION_BATCH_SIZE, + ); + removed += batch; + // A short batch means the window is clean; the next sweep picks up + // whatever aged into it meanwhile. + if (batch < RETENTION_BATCH_SIZE) break; + } + return removed; + } + /** Mark a notification as shown (user saw it) and push the ack event. */ async markShown(uid: string, userId: number): Promise { await this.stores.notification.markShown(uid, userId); @@ -266,6 +308,22 @@ export class NotificationService extends PuterService { // -- Internals --------------------------------------------------- + #retentionDays(): number { + const configured = Number(this.config.notificationRetentionDays ?? 0); + return Number.isFinite(configured) && configured > 0 ? configured : 0; + } + + #armRetentionSweep(): void { + if (this.#retentionDays() <= 0) return; + const sweep = setInterval(() => { + void this.sweepExpired().catch((err) => { + console.warn('[notification] retention sweep failed', err); + }); + }, RETENTION_SWEEP_INTERVAL_MS); + sweep.unref?.(); + this.#retentionSweep = sweep; + } + async #sendUnreads(userId: number): Promise { // Fetch all unseen + unacknowledged notifications const rows = await this.stores.notification.listByUserId(userId, { diff --git a/src/backend/stores/notification/NotificationStore.js b/src/backend/stores/notification/NotificationStore.js index adb0aecfd..92e3f9e1d 100644 --- a/src/backend/stores/notification/NotificationStore.js +++ b/src/backend/stores/notification/NotificationStore.js @@ -185,6 +185,41 @@ export class NotificationStore extends PuterStore { return changed; } + /** + * Delete up to `limit` rows created more than `days` ago, and report how + * many went — a full batch means there is more behind it. + * + * Only mysql takes a LIMIT on DELETE — postgres has none and sqlite's needs + * an optional build flag — so the other two bound the batch through an id + * list. Each engine computes its own cutoff, so no clock crosses the wire. + * + * @param {number} days @param {number} limit + */ + async deleteCreatedBefore(days, limit) { + const retentionDays = Math.floor(Number(days)); + const batch = Math.floor(Number(limit)); + if (!Number.isFinite(retentionDays) || retentionDays <= 0) return 0; + if (!Number.isFinite(batch) || batch <= 0) return 0; + + const cutoff = this.clients.db.case({ + sqlite: `datetime('now', '-${retentionDays} days')`, + postgres: `(NOW() - INTERVAL '${retentionDays} days')`, + otherwise: `(NOW() - INTERVAL ${retentionDays} DAY)`, + }); + const statement = this.clients.db.case({ + mysql: + 'DELETE FROM `notification` ' + + `WHERE \`created_at\` < ${cutoff} ORDER BY \`id\` LIMIT ?`, + otherwise: + 'DELETE FROM `notification` WHERE `id` IN (' + + 'SELECT `id` FROM `notification` ' + + `WHERE \`created_at\` < ${cutoff} ORDER BY \`id\` LIMIT ?)`, + }); + + const result = await this.clients.db.write(statement, [batch]); + return result?.affectedRows ?? result?.changes ?? 0; + } + // -- Internals ---------------------------------------------------- #unackCacheKey(userId) { diff --git a/src/backend/stores/notification/NotificationStore.test.js b/src/backend/stores/notification/NotificationStore.test.js index 687322a83..b2fc8bb6c 100644 --- a/src/backend/stores/notification/NotificationStore.test.js +++ b/src/backend/stores/notification/NotificationStore.test.js @@ -322,6 +322,90 @@ describe('NotificationStore', () => { expect(await redis.get(unackKey(u.id))).toBe('1'); }); + // -- retention ----------------------------------------------------- + + /** + * Age a row by rewriting `created_at`. The format is what every engine + * writes for a timestamp column, so the comparison the sweep makes is the + * one production makes. + */ + const backdate = async (uid, days) => { + const when = new Date(Date.now() - days * 86_400_000) + .toISOString() + .replace('T', ' ') + .slice(0, 19); + await server.clients.db.write( + 'UPDATE `notification` SET `created_at` = ? WHERE `uid` = ?', + [when, uid], + ); + }; + + /** Clear anything an earlier test aged, so counts below are exact. */ + const drain = async () => { + while ((await store.deleteCreatedBefore(14, 500)) > 0); + }; + + it('deletes rows past the window and leaves the ones inside it', async () => { + await drain(); + const u = await makeUser(); + const old = await store.create({ userId: u.id, value: { n: 'old' } }); + const alsoOld = await store.create({ userId: u.id, value: { n: '2' } }); + const recent = await store.create({ userId: u.id, value: { n: 'new' } }); + await backdate(old.uid, 20); + await backdate(alsoOld.uid, 15); + await backdate(recent.uid, 13); + + expect(await store.deleteCreatedBefore(14, 500)).toBe(2); + expect(await store.getByUid(old.uid)).toBeNull(); + expect(await store.getByUid(alsoOld.uid)).toBeNull(); + expect((await store.getByUid(recent.uid))?.uid).toBe(recent.uid); + }); + + it('takes acknowledged rows and unacknowledged ones alike', async () => { + await drain(); + const u = await makeUser(); + const acked = await store.create({ userId: u.id, value: {} }); + const never = await store.create({ userId: u.id, value: {} }); + await store.markAcknowledged(acked.uid, u.id); + await backdate(acked.uid, 20); + await backdate(never.uid, 20); + + expect(await store.deleteCreatedBefore(14, 500)).toBe(2); + expect(await store.listByUserId(u.id)).toEqual([]); + }); + + it('stops at the batch size so the caller can keep going', async () => { + await drain(); + const u = await makeUser(); + for (let i = 0; i < 5; i++) { + const row = await store.create({ userId: u.id, value: { i } }); + await backdate(row.uid, 20); + } + + expect(await store.deleteCreatedBefore(14, 2)).toBe(2); + expect(await store.deleteCreatedBefore(14, 2)).toBe(2); + expect(await store.deleteCreatedBefore(14, 2)).toBe(1); + expect(await store.deleteCreatedBefore(14, 2)).toBe(0); + }); + + it('deletes nothing for a window or batch that is not a positive count', async () => { + await drain(); + const u = await makeUser(); + const n = await store.create({ userId: u.id, value: {} }); + await backdate(n.uid, 40); + + for (const [days, limit] of [ + [0, 500], + [-1, 500], + ['forever', 500], + [14, 0], + [14, -5], + ]) { + expect(await store.deleteCreatedBefore(days, limit)).toBe(0); + } + expect((await store.getByUid(n.uid))?.uid).toBe(n.uid); + }); + it('will not delete another user notification', async () => { const u = await makeUser(); const n = await store.create({ userId: u.id, value: {} }); diff --git a/src/backend/types.ts b/src/backend/types.ts index 41ccede91..f37aa1062 100644 --- a/src/backend/types.ts +++ b/src/backend/types.ts @@ -911,6 +911,13 @@ interface IConfigOptional { share_max_recipients?: number; share_max_items?: number; + /** + * How long a notification is kept, in days from creation — the mailbox is + * not an archive, so acknowledged or not, a row past this goes. Omit it, or + * set 0, and nothing is ever swept. + */ + notificationRetentionDays?: number; + // -- Storage / S3 ------------------------------------------------ /** S3 storage config (local fauxqs or remote). */