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
14 changes: 14 additions & 0 deletions backend/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,20 @@ INDEXER_REORG_ROLLBACK_DEPTH=10
LEADERBOARD_SNAPSHOT_CRON=0 * * * *
LEADERBOARD_SNAPSHOT_RETENTION_DAYS=30

# Account Data Export Configuration
EXPORT_DIR=./exports
# Hours a ready export stays downloadable before its file is cleaned up
EXPORT_TTL_HOURS=48
# Stale export file cleanup schedule
EXPORT_CLEANUP_ENABLED=true
EXPORT_CLEANUP_CRON=0 * * * *
# Hours a failed export job row is kept before removal
EXPORT_FAILED_RETENTION_HOURS=24
# Minutes a job may stay in "processing" before it is treated as stuck and failed
EXPORT_STUCK_PROCESSING_MINUTES=60
# Minimum age (minutes) of an export file with no job row before it is deleted
EXPORT_ORPHAN_GRACE_MINUTES=60

# Idempotency Key Configuration
IDEMPOTENCY_KEY_TTL_HOURS=24

Expand Down
89 changes: 89 additions & 0 deletions backend/src/account/account-export-cleanup.scheduler.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
import {
Injectable,
Logger,
OnModuleDestroy,
OnModuleInit,
} from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import { SchedulerRegistry } from '@nestjs/schedule';
import { CronJob } from 'cron';
import { AccountService } from './account.service';

export const EXPORT_CLEANUP_JOB_NAME = 'account-export-stale-file-cleanup';

/**
* Schedules periodic removal of stale account export files and job rows on
* the cadence given by EXPORT_CLEANUP_CRON. Runs never overlap: a tick that
* fires while the previous run is still in progress is skipped.
*/
@Injectable()
export class AccountExportCleanupScheduler
implements OnModuleInit, OnModuleDestroy
{
private readonly logger = new Logger(AccountExportCleanupScheduler.name);
private running = false;

constructor(
private readonly accountService: AccountService,
private readonly schedulerRegistry: SchedulerRegistry,
private readonly configService: ConfigService,
) {}

onModuleInit(): void {
const enabled =
String(
this.configService.get<string>('EXPORT_CLEANUP_ENABLED', 'true'),
).toLowerCase() !== 'false';
if (!enabled) {
this.logger.log('Account export cleanup is disabled');
return;
}

const cronExpression = this.configService.get<string>(
'EXPORT_CLEANUP_CRON',
'0 * * * *',
);

const job = new CronJob(cronExpression, () => {
void this.handleCleanup();
});

this.schedulerRegistry.addCronJob(EXPORT_CLEANUP_JOB_NAME, job);
job.start();

this.logger.log(
`Account export cleanup scheduled with cron "${cronExpression}"`,
);
}

onModuleDestroy(): void {
if (this.schedulerRegistry.doesExist('cron', EXPORT_CLEANUP_JOB_NAME)) {
this.schedulerRegistry.deleteCronJob(EXPORT_CLEANUP_JOB_NAME);
}
}

async handleCleanup(): Promise<void> {
if (this.running) {
this.logger.warn(
'Previous account export cleanup still running; skipping this tick',
);
return;
}

this.running = true;
try {
const { expired, failed, stuck, orphans } =
await this.accountService.cleanupExports();
if (expired + failed + stuck + orphans > 0) {
this.logger.log(
`Account export cleanup: expired=${expired} failed=${failed} ` +
`stuck=${stuck} orphans=${orphans}`,
);
}
} catch (err) {
this.logger.error('Account export cleanup failed', err);
} finally {
this.running = false;
}
}
}
3 changes: 2 additions & 1 deletion backend/src/account/account.module.ts
Original file line number Diff line number Diff line change
@@ -1,12 +1,13 @@
import { Module } from '@nestjs/common';
import { TypeOrmModule } from '@nestjs/typeorm';
import { AccountExportCleanupScheduler } from './account-export-cleanup.scheduler';
import { AccountController } from './account.controller';
import { AccountService } from './account.service';
import { DataExportJob } from './entities/data-export-job.entity';

@Module({
imports: [TypeOrmModule.forFeature([DataExportJob])],
controllers: [AccountController],
providers: [AccountService],
providers: [AccountService, AccountExportCleanupScheduler],
})
export class AccountModule {}
137 changes: 132 additions & 5 deletions backend/src/account/account.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,11 +7,21 @@ import {
import { ConfigService } from '@nestjs/config';
import { InjectDataSource, InjectRepository } from '@nestjs/typeorm';
import { Cron, CronExpression } from '@nestjs/schedule';
import { DataSource, LessThan, Repository } from 'typeorm';
import { DataSource, In, LessThan, Repository } from 'typeorm';
import * as fs from 'fs/promises';
import * as path from 'path';
import { DataExportJob } from './entities/data-export-job.entity';

const UUID_PATTERN =
/^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i;

export interface ExportCleanupResult {
expired: number;
failed: number;
stuck: number;
orphans: number;
}

@Injectable()
export class AccountService {
constructor(
Expand Down Expand Up @@ -241,14 +251,131 @@ export class AccountService {
}
}

@Cron(CronExpression.EVERY_DAY_AT_MIDNIGHT)
async cleanupExports(): Promise<void> {
/**
* Removes stale export artifacts. Invoked on a configurable schedule by
* {@link AccountExportCleanupScheduler}.
*
* - `ready` jobs past `expires_at`: file unlinked, row deleted
* - `failed` jobs older than EXPORT_FAILED_RETENTION_HOURS: file unlinked, row deleted
* - `processing` jobs older than EXPORT_STUCK_PROCESSING_MINUTES: partial file
* unlinked, job marked `failed` (deleted on a later run by the rule above)
* - files in EXPORT_DIR with no matching job row and older than
* EXPORT_ORPHAN_GRACE_MINUTES: unlinked
*/
async cleanupExports(): Promise<ExportCleanupResult> {
const now = Date.now();
const result: ExportCleanupResult = {
expired: 0,
failed: 0,
stuck: 0,
orphans: 0,
};

const expired = await this.jobRepo.find({
where: { status: 'ready', expires_at: LessThan(new Date()) },
where: { status: 'ready', expires_at: LessThan(new Date(now)) },
});
for (const job of expired) {
if (job.file_path) await fs.unlink(job.file_path).catch(() => {});
await this.removeJobFile(job);
await this.jobRepo.delete(job.id);
result.expired++;
}

const failedRetentionHours = Number(
this.configService.get<number>('EXPORT_FAILED_RETENTION_HOURS', 24),
);
const failed = await this.jobRepo.find({
where: {
status: 'failed',
updated_at: LessThan(new Date(now - failedRetentionHours * 3_600_000)),
},
});
for (const job of failed) {
await this.removeJobFile(job);
await this.jobRepo.delete(job.id);
result.failed++;
}

const stuckMinutes = Number(
this.configService.get<number>('EXPORT_STUCK_PROCESSING_MINUTES', 60),
);
const stuck = await this.jobRepo.find({
where: {
status: 'processing',
updated_at: LessThan(new Date(now - stuckMinutes * 60_000)),
},
});
for (const job of stuck) {
await this.removeJobFile(job);
await this.jobRepo.update(job.id, {
status: 'failed',
file_path: null,
expires_at: null,
});
result.stuck++;
}

result.orphans = await this.removeOrphanFiles(now);

return result;
}

private exportDir(): string {
return this.configService.get<string>('EXPORT_DIR', './exports');
}

private async removeJobFile(job: DataExportJob): Promise<void> {
const filePath =
job.file_path ?? path.join(this.exportDir(), `${job.id}.json`);
await fs.unlink(filePath).catch(() => {});
}

private async removeOrphanFiles(now: number): Promise<number> {
const dir = this.exportDir();
let entries: string[];
try {
entries = await fs.readdir(dir);
} catch {
return 0;
}

const graceMinutes = Number(
this.configService.get<number>('EXPORT_ORPHAN_GRACE_MINUTES', 60),
);
const cutoff = now - graceMinutes * 60_000;

// Only consider files named like export jobs (`<uuid>.json`) so unrelated
// files that happen to live in EXPORT_DIR are never touched.
const candidates: { jobId: string; filePath: string }[] = [];
for (const name of entries) {
const jobId = path.basename(name, '.json');
if (!name.endsWith('.json') || !UUID_PATTERN.test(jobId)) continue;
const filePath = path.join(dir, name);
try {
const stat = await fs.stat(filePath);
if (!stat.isFile() || stat.mtimeMs > cutoff) continue;
} catch {
continue;
}
candidates.push({ jobId, filePath });
}
if (candidates.length === 0) return 0;

const knownJobs = await this.jobRepo.find({
where: { id: In(candidates.map((c) => c.jobId)) },
select: ['id'],
});
const known = new Set(knownJobs.map((j) => j.id));

let removed = 0;
for (const c of candidates) {
if (known.has(c.jobId)) continue;
try {
await fs.unlink(c.filePath);
removed++;
} catch {
// already gone or not removable; skip
}
}
return removed;
}
}
Loading