mirror of
https://github.com/snapotter-hq/SnapOtter.git
synced 2026-08-03 07:46:42 +02:00
feat(enterprise): add per-team retention overrides with deleteAfter
Compute a deleteAfter timestamp on job creation when the enterprise team_retention_overrides feature is enabled. The cleanup sweep now deletes storage for jobs past their deleteAfter deadline, running independently of the global TTL setting.
This commit is contained in:
@@ -6,6 +6,7 @@
|
||||
* until the worker produces a result or the sync-wait window expires.
|
||||
*/
|
||||
import { FlowProducer, type Job, QueueEvents } from "bullmq";
|
||||
import { eq } from "drizzle-orm";
|
||||
import { env } from "../config.js";
|
||||
import { db, schema } from "../db/index.js";
|
||||
import { createRedisConnection } from "./connection.js";
|
||||
@@ -76,6 +77,11 @@ export async function enqueueToolJob(data: ToolJobData): Promise<Job<ToolJobData
|
||||
settings: (data.dbSettings ?? data.settings) as Record<string, unknown>,
|
||||
});
|
||||
|
||||
// Fire-and-forget: compute deleteAfter from team retention override
|
||||
if (data.userId) {
|
||||
void computeDeleteAfter(data.jobId, data.userId).catch(() => {});
|
||||
}
|
||||
|
||||
const queue = getQueue(data.pool);
|
||||
const job = await queue.add(data.toolId, { ...data, jobId: data.jobId }, { jobId: data.jobId });
|
||||
return job;
|
||||
@@ -108,3 +114,43 @@ export async function waitForJob(
|
||||
throw err; // real failure
|
||||
}
|
||||
}
|
||||
|
||||
// ── Per-team retention override ────────────────────────────────
|
||||
|
||||
/**
|
||||
* Compute and set `deleteAfter` on a job row based on the owning user's
|
||||
* team retention setting. Only applies when the enterprise
|
||||
* `team_retention_overrides` feature is enabled. Fire-and-forget; failures
|
||||
* never block job creation.
|
||||
*/
|
||||
async function computeDeleteAfter(jobId: string, userId: string): Promise<void> {
|
||||
let isTeamRetentionEnabled = false;
|
||||
try {
|
||||
const { isFeatureEnabled } = await import("@snapotter/enterprise");
|
||||
isTeamRetentionEnabled = isFeatureEnabled("team_retention_overrides");
|
||||
} catch {}
|
||||
|
||||
if (!isTeamRetentionEnabled) return;
|
||||
|
||||
const userRow = await db
|
||||
.select({ team: schema.users.team })
|
||||
.from(schema.users)
|
||||
.where(eq(schema.users.id, userId))
|
||||
.limit(1);
|
||||
|
||||
if (!userRow.length || !userRow[0].team) return;
|
||||
|
||||
const teamRow = await db
|
||||
.select({ retentionHours: schema.teams.retentionHours })
|
||||
.from(schema.teams)
|
||||
.where(eq(schema.teams.id, userRow[0].team))
|
||||
.limit(1);
|
||||
|
||||
const retentionHours =
|
||||
teamRow.length && teamRow[0].retentionHours !== null
|
||||
? teamRow[0].retentionHours
|
||||
: env.FILE_MAX_AGE_HOURS;
|
||||
|
||||
const deleteAfter = new Date(Date.now() + retentionHours * 60 * 60 * 1000);
|
||||
await db.update(schema.jobs).set({ deleteAfter }).where(eq(schema.jobs.id, jobId));
|
||||
}
|
||||
|
||||
@@ -10,7 +10,7 @@
|
||||
* calling runSystemJob); anything else is a bug.
|
||||
*/
|
||||
import type { Job } from "bullmq";
|
||||
import { eq, inArray, sql } from "drizzle-orm";
|
||||
import { and, eq, inArray, isNotNull, lt, sql } from "drizzle-orm";
|
||||
import { env } from "../config.js";
|
||||
import { db, schema } from "../db/index.js";
|
||||
import { getMaxAgeMs } from "../lib/cleanup.js";
|
||||
@@ -121,8 +121,35 @@ export function decideExpiry(
|
||||
}
|
||||
|
||||
async function storageTtlSweep(): Promise<{ removed: number; failed: number }> {
|
||||
// --- Per-job deleteAfter sweep (team retention overrides) ---
|
||||
// Runs regardless of the global TTL; deleteAfter is an absolute deadline.
|
||||
let deleteAfterCleaned = 0;
|
||||
try {
|
||||
const expiredJobs = await db
|
||||
.select({ id: schema.jobs.id })
|
||||
.from(schema.jobs)
|
||||
.where(and(isNotNull(schema.jobs.deleteAfter), lt(schema.jobs.deleteAfter, new Date())));
|
||||
|
||||
for (const job of expiredJobs) {
|
||||
try {
|
||||
await deletePrefix(`uploads/${job.id}`);
|
||||
await deletePrefix(`outputs/${job.id}`);
|
||||
deleteAfterCleaned++;
|
||||
} catch {
|
||||
// Directory may not exist
|
||||
}
|
||||
}
|
||||
|
||||
if (deleteAfterCleaned > 0) {
|
||||
console.log(`Storage TTL: cleaned up ${deleteAfterCleaned} jobs by deleteAfter`);
|
||||
}
|
||||
} catch {
|
||||
// deleteAfter sweep is best-effort
|
||||
}
|
||||
|
||||
// --- Global TTL sweep ---
|
||||
const maxAgeMs = await getMaxAgeMs();
|
||||
if (maxAgeMs <= 0) return { removed: 0, failed: 0 };
|
||||
if (maxAgeMs <= 0) return { removed: deleteAfterCleaned, failed: 0 };
|
||||
|
||||
const cutoffMs = Date.now() - maxAgeMs;
|
||||
const uploadDirs = await listJobDirs("uploads");
|
||||
@@ -168,7 +195,7 @@ async function storageTtlSweep(): Promise<{ removed: number; failed: number }> {
|
||||
if (removed > 0) {
|
||||
console.log(`Storage TTL: removed ${removed} expired job dirs`);
|
||||
}
|
||||
return { removed, failed: errors.length };
|
||||
return { removed: removed + deleteAfterCleaned, failed: errors.length };
|
||||
}
|
||||
|
||||
// -- Retention sweep ----------------------------------------------------------
|
||||
|
||||
Reference in New Issue
Block a user