Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
87 changes: 87 additions & 0 deletions src/services/DeletionPurgeJob.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
import { prisma } from 'src/db';
import { logAuditEvent } from 'src/utils/auditLogger';

/**
* DeletionPurgeJob — executes pending account deletions whose grace period has elapsed.
*
* Called by the nightly deletion-purge scheduler. The job is idempotent:
* if a user record has already been removed (e.g. manual admin action) the
* corresponding stale pendingDeletion row is cleaned up without error.
*
* Cascade behaviour (from schema.prisma):
* - pendingDeletion.onDelete = Cascade → removed with the User row
* - auditLog.onDelete = SetNull → audit trail is preserved (userId → NULL)
*/
export class DeletionPurgeJob {
/**
* Run the purge sweep for all accounts whose scheduledAt has passed.
* Call this from a nightly cron/scheduler.
*/
static async run(): Promise<PurgeRunSummary> {
const summary: PurgeRunSummary = { purged: 0, skipped: 0, errors: [] };

const due = await prisma.pendingDeletion.findMany({
where: { scheduledAt: { lte: new Date() } },
select: { userId: true },
});

for (const { userId } of due) {
try {
const deleted = await this.purgeUser(userId);
if (deleted) {
summary.purged++;
} else {
summary.skipped++;
}
} catch (error) {
summary.errors.push({
userId,
message: error instanceof Error ? error.message : String(error),
});
}
}

return summary;
}

/**
* Purge a single user account.
*
* Returns true if the user was deleted, false if the user was already gone
* (stale record cleaned up instead). Exported so tests and admin tooling
* can trigger an immediate per-user purge.
*/
static async purgeUser(userId: string): Promise<boolean> {
const user = await prisma.user.findUnique({
where: { id: userId },
select: { id: true },
});

if (!user) {
// User already removed — clean up the orphaned pendingDeletion row.
await prisma.pendingDeletion.deleteMany({ where: { userId } });
return false;
}

// Write audit record before deletion — the AuditLog relation uses
// onDelete: SetNull so this row survives the cascade as a compliance trail.
await logAuditEvent(userId, 'DATA_DELETE', {
entityType: 'User',
entityId: userId,
statusCode: 200,
metadata: { action: 'purge_account' },
});

// Deleting the User row cascades to pendingDeletion (and all other
// Cascade-linked relations) automatically.
await prisma.user.delete({ where: { id: userId } });

return true;
}
}

export interface PurgeRunSummary {
purged: number;
skipped: number;
errors: Array<{ userId: string; message: string }>;
}
180 changes: 180 additions & 0 deletions src/services/__tests__/DeletionPurgeJob.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,180 @@
import { afterAll, afterEach, beforeAll, describe, expect, it } from 'vitest';

import { prisma } from 'src/db';
import { DeletionPurgeJob } from 'src/services/DeletionPurgeJob';

function daysAgo(days: number): Date {
return new Date(Date.now() - days * 24 * 60 * 60 * 1000);
}

function daysFromNow(days: number): Date {
return new Date(Date.now() + days * 24 * 60 * 60 * 1000);
}

async function createUser(tag: string) {
const runId = Date.now();
return prisma.user.create({
data: {
email: `purge-${tag}-${runId}@example.com`,
password: 'hashed',
firstName: 'Purge',
lastName: tag,
},
});
}

describe('DeletionPurgeJob', () => {
// Users that must survive all tests — cleaned up in afterAll
let spectatorId: string;

beforeAll(async () => {
const spectator = await createUser('spectator');
spectatorId = spectator.id;
});

afterEach(async () => {
// Remove any pendingDeletion rows left over from failed tests so the
// next test starts clean. User rows created per-test are cleaned up
// inside each test's own try/finally block.
await prisma.pendingDeletion.deleteMany({
where: { scheduledAt: { lte: daysFromNow(1) } },
});
});

afterAll(async () => {
await prisma.user.deleteMany({ where: { id: spectatorId } });
});

// -----------------------------------------------------------------------
// purgeUser — per-user behaviour
// -----------------------------------------------------------------------

it('purgeUser() deletes a user and returns true', async () => {
const user = await createUser('to-delete');

await prisma.pendingDeletion.create({
data: { userId: user.id, scheduledAt: daysAgo(1) },
});

const deleted = await DeletionPurgeJob.purgeUser(user.id);

expect(deleted).toBe(true);
const found = await prisma.user.findUnique({ where: { id: user.id } });
expect(found).toBeNull();
});

it('purgeUser() is idempotent when the user no longer exists', async () => {
const user = await createUser('already-gone');
const userId = user.id;

// Create a pendingDeletion row then manually delete the user to simulate
// a race / prior manual deletion.
await prisma.pendingDeletion.create({
data: { userId, scheduledAt: daysAgo(1) },
});
await prisma.user.delete({ where: { id: userId } });

// At this point the pendingDeletion row may or may not have been
// cascade-deleted; either way purgeUser should not throw.
const deleted = await DeletionPurgeJob.purgeUser(userId);

expect(deleted).toBe(false);
});

it('purgeUser() writes an audit record before deleting the user', async () => {
const user = await createUser('audited');

await prisma.pendingDeletion.create({
data: { userId: user.id, scheduledAt: daysAgo(1) },
});

const before = new Date();
await DeletionPurgeJob.purgeUser(user.id);
const after = new Date();

// AuditLog uses onDelete: SetNull, so the row survives user deletion.
const logs = await prisma.auditLog.findMany({
where: {
entityId: user.id,
action: 'DATA_DELETE',
createdAt: { gte: before, lte: after },
},
});

expect(logs.length).toBeGreaterThanOrEqual(1);

// Cleanup orphaned audit rows
await prisma.auditLog.deleteMany({ where: { entityId: user.id } });
});

// -----------------------------------------------------------------------
// run() — batch sweep
// -----------------------------------------------------------------------

it('run() purges due deletions and leaves future ones untouched', async () => {
const dueUser = await createUser('due');
const futureUser = await createUser('future');

await prisma.pendingDeletion.createMany({
data: [
{ userId: dueUser.id, scheduledAt: daysAgo(1) },
{ userId: futureUser.id, scheduledAt: daysFromNow(30) },
],
});

try {
const summary = await DeletionPurgeJob.run();

// dueUser must be gone
const dueFound = await prisma.user.findUnique({ where: { id: dueUser.id } });
expect(dueFound).toBeNull();

// futureUser must still exist
const futureFound = await prisma.user.findUnique({ where: { id: futureUser.id } });
expect(futureFound).not.toBeNull();

expect(summary.purged).toBeGreaterThanOrEqual(1);
const ourErrors = summary.errors.filter(
(e) => e.userId === dueUser.id || e.userId === futureUser.id,
);
expect(ourErrors).toHaveLength(0);
} finally {
await prisma.pendingDeletion.deleteMany({ where: { userId: futureUser.id } });
await prisma.user.deleteMany({ where: { id: futureUser.id } });
// dueUser was already deleted by the job
}
});

it('run() counts a stale record (user pre-deleted) as skipped, not an error', async () => {
const user = await createUser('stale');
const userId = user.id;

await prisma.pendingDeletion.create({
data: { userId, scheduledAt: daysAgo(1) },
});

// Simulate prior manual deletion — cascade may already remove pendingDeletion;
// if not, we want run() to handle the stale row gracefully.
try {
await prisma.user.delete({ where: { id: userId } });
} catch {
// Already gone — fine
}

const summary = await DeletionPurgeJob.run();

const ourErrors = summary.errors.filter((e) => e.userId === userId);
expect(ourErrors).toHaveLength(0);
});

it('run() does not touch users without a due pendingDeletion', async () => {
// spectatorId has no pendingDeletion row at all
const summary = await DeletionPurgeJob.run();

const spectator = await prisma.user.findUnique({ where: { id: spectatorId } });
expect(spectator).not.toBeNull();

const ourErrors = summary.errors.filter((e) => e.userId === spectatorId);
expect(ourErrors).toHaveLength(0);
});
});
52 changes: 52 additions & 0 deletions src/utils/deletionPurgeScheduler.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
import { DeletionPurgeJob } from 'src/services/DeletionPurgeJob';
import { env } from './helpers';

const LOG_PREFIX = '[Deletion Purge Scheduler]';

const DAY = 24 * 60 * 60 * 1000;

let purgeInterval: NodeJS.Timeout | null = null;

const runPurge = async () => {
try {
if (env('NODE_ENV') !== 'test') {
console.log(`${LOG_PREFIX} Running deletion purge sweep...`);
}

const summary = await DeletionPurgeJob.run();

if (env('NODE_ENV') !== 'test') {
console.log(
`${LOG_PREFIX} Sweep complete — purged: ${summary.purged}, skipped: ${summary.skipped}, errors: ${summary.errors.length}`,
);

for (const { userId, message } of summary.errors) {
console.error(`${LOG_PREFIX} Error purging user ${userId}: ${message}`);
}
}
} catch (error) {
console.error(`${LOG_PREFIX} Sweep failed:`, error);
}
};

export const startDeletionPurgeScheduler = () => {
if (env('NODE_ENV') === 'test') {
return;
}

console.log(`${LOG_PREFIX} Starting scheduled deletion purge job (every 24 h)...`);

purgeInterval = setInterval(runPurge, DAY);
};

export const stopDeletionPurgeScheduler = () => {
if (purgeInterval) {
clearInterval(purgeInterval);
purgeInterval = null;
}
};

export default {
start: startDeletionPurgeScheduler,
stop: stopDeletionPurgeScheduler,
};
2 changes: 2 additions & 0 deletions src/utils/initialize.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ import methodOverride from "method-override";
import passport from "passport";
import path from "path";
import { startAnalyticsScheduler } from "./analyticsScheduler";
import { startDeletionPurgeScheduler } from "./deletionPurgeScheduler";
import { startMediaScheduler } from "./mediaScheduler";
import { startMonitoringScheduler } from "src/services/monitoringService";

Expand Down Expand Up @@ -156,6 +157,7 @@ export const initialize = async (app: Express) => {
startLogCleanupScheduler();
startAnalyticsScheduler();
startMediaScheduler();
startDeletionPurgeScheduler();

if (process.env.NODE_ENV !== "test") {
console.log("[Security] All security services initialized successfully");
Expand Down
Loading