feat: complete TASK-WP2-03 project lifecycle
This commit is contained in:
@@ -0,0 +1,241 @@
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { rmSync } from "node:fs";
|
||||
import { createRequire } from "node:module";
|
||||
import { isAbsolute, relative, resolve } from "node:path";
|
||||
|
||||
import type BetterSqlite3 from "better-sqlite3";
|
||||
|
||||
const require = createRequire(import.meta.url);
|
||||
const Database = require("better-sqlite3") as typeof BetterSqlite3;
|
||||
|
||||
interface ProjectCleanupRow {
|
||||
cleanup_id: string;
|
||||
project_id: string;
|
||||
}
|
||||
|
||||
interface ExpiredProjectRow {
|
||||
owner_id: string;
|
||||
project_id: string;
|
||||
}
|
||||
|
||||
interface FileCleanupRow {
|
||||
byte_size: number;
|
||||
cleanup_id: string;
|
||||
counts_toward_managed: 0 | 1;
|
||||
managed_file_id: string | null;
|
||||
relative_path: string;
|
||||
}
|
||||
|
||||
const resourceScope = JSON.stringify([
|
||||
"project_state", "generation", "generated_image", "reference", "location", "latest_export",
|
||||
]);
|
||||
|
||||
function iso(timestamp: number) {
|
||||
return new Date(timestamp).toISOString();
|
||||
}
|
||||
|
||||
export class ProjectPurgeCleanup {
|
||||
private readonly clock: () => number;
|
||||
private readonly dataRoot: string | undefined;
|
||||
private readonly database: BetterSqlite3.Database;
|
||||
|
||||
constructor(input: { clock?: () => number; dataRoot?: string; databasePath: string }) {
|
||||
this.clock = input.clock ?? Date.now;
|
||||
this.dataRoot = input.dataRoot ? resolve(input.dataRoot) : undefined;
|
||||
const nativeBinding = process.env.DADA_SQLITE_NATIVE_BINDING;
|
||||
this.database = new Database(input.databasePath, nativeBinding ? { nativeBinding } : undefined);
|
||||
this.database.pragma("journal_mode = WAL");
|
||||
this.database.pragma("foreign_keys = ON");
|
||||
this.database.pragma("busy_timeout = 5000");
|
||||
}
|
||||
|
||||
close() {
|
||||
this.database.close();
|
||||
}
|
||||
|
||||
run() {
|
||||
const expired = this.sweepExpired();
|
||||
const relational = this.processPending();
|
||||
const physical = this.processFileCleanup();
|
||||
return { expired, physical, relational };
|
||||
}
|
||||
|
||||
sweepExpired() {
|
||||
if (!this.tableExists("projects") || !this.tableExists("project_cleanup_queue")) return 0;
|
||||
const now = this.clock();
|
||||
const rows = this.database.prepare(`
|
||||
SELECT project_id, owner_id FROM projects
|
||||
WHERE status = 'trashed' AND purge_at <= ?
|
||||
ORDER BY purge_at, project_id
|
||||
`).all(now) as ExpiredProjectRow[];
|
||||
let expired = 0;
|
||||
for (const row of rows) {
|
||||
const transaction = this.database.transaction(() => {
|
||||
const changed = this.database.prepare(`
|
||||
UPDATE projects SET status = 'purged', updated_at = ?, state_version = state_version + 1
|
||||
WHERE project_id = ? AND status = 'trashed'
|
||||
`).run(now, row.project_id);
|
||||
if (changed.changes !== 1) return false;
|
||||
this.queueManagedFiles(row.project_id, now);
|
||||
this.database.prepare(`
|
||||
INSERT OR IGNORE INTO project_cleanup_queue (
|
||||
cleanup_id, project_id, owner_id, resource_scope_json, status, created_at, completed_at, last_error
|
||||
) VALUES (?, ?, ?, ?, 'pending', ?, NULL, NULL)
|
||||
`).run(randomUUID(), row.project_id, row.owner_id, resourceScope, now);
|
||||
return true;
|
||||
});
|
||||
if (transaction.immediate()) expired += 1;
|
||||
}
|
||||
return expired;
|
||||
}
|
||||
|
||||
processPending() {
|
||||
if (!this.tableExists("project_cleanup_queue") || !this.tableExists("projects")) return { completed: 0, failed: 0 };
|
||||
const rows = this.database.prepare(`
|
||||
SELECT cleanup_id, project_id FROM project_cleanup_queue
|
||||
WHERE status IN ('pending', 'failed') ORDER BY created_at, project_id
|
||||
`).all() as ProjectCleanupRow[];
|
||||
let completed = 0;
|
||||
let failed = 0;
|
||||
for (const row of rows) {
|
||||
try {
|
||||
const transaction = this.database.transaction(() => {
|
||||
this.database.prepare("DELETE FROM projects WHERE project_id = ? AND status = 'purged'").run(row.project_id);
|
||||
this.database.prepare(`
|
||||
UPDATE project_cleanup_queue
|
||||
SET status = 'completed', completed_at = ?, last_error = NULL
|
||||
WHERE cleanup_id = ?
|
||||
`).run(this.clock(), row.cleanup_id);
|
||||
});
|
||||
transaction.immediate();
|
||||
completed += 1;
|
||||
} catch {
|
||||
this.database.prepare(`
|
||||
UPDATE project_cleanup_queue SET status = 'failed', last_error = 'project_cleanup_failed'
|
||||
WHERE cleanup_id = ?
|
||||
`).run(row.cleanup_id);
|
||||
failed += 1;
|
||||
}
|
||||
}
|
||||
return { completed, failed };
|
||||
}
|
||||
|
||||
processFileCleanup() {
|
||||
if (!this.dataRoot || !this.tableExists("file_cleanup_queue") || !this.tableExists("managed_files")) {
|
||||
return { completed: 0, failed: 0 };
|
||||
}
|
||||
const rows = this.database.prepare(`
|
||||
SELECT cleanup_id, managed_file_id, relative_path, byte_size, counts_toward_managed
|
||||
FROM file_cleanup_queue WHERE status IN ('pending', 'failed') ORDER BY created_at, cleanup_id
|
||||
`).all() as FileCleanupRow[];
|
||||
let completed = 0;
|
||||
let failed = 0;
|
||||
for (const row of rows) {
|
||||
try {
|
||||
const path = this.resolveManagedPath(row.relative_path);
|
||||
rmSync(path, { force: true });
|
||||
const transaction = this.database.transaction(() => {
|
||||
if (row.counts_toward_managed === 1 && row.managed_file_id) {
|
||||
if (this.tableExists("project_asset_refs")) {
|
||||
this.database.prepare("DELETE FROM project_asset_refs WHERE managed_file_id = ?").run(row.managed_file_id);
|
||||
}
|
||||
if (this.tableExists("project_resource_files")) {
|
||||
this.database.prepare("DELETE FROM project_resource_files WHERE managed_file_id = ?").run(row.managed_file_id);
|
||||
}
|
||||
if (this.tableExists("asset_cleanup_request_items")) {
|
||||
const requests = this.database.prepare("SELECT request_id FROM asset_cleanup_request_items WHERE managed_file_id = ?")
|
||||
.all(row.managed_file_id) as Array<{ request_id: string }>;
|
||||
this.database.prepare("DELETE FROM asset_cleanup_request_items WHERE managed_file_id = ?").run(row.managed_file_id);
|
||||
if (this.tableExists("asset_cleanup_requests")) {
|
||||
for (const request of requests) {
|
||||
const pending = this.database.prepare("SELECT COUNT(*) AS count FROM asset_cleanup_request_items WHERE request_id = ?")
|
||||
.get(request.request_id) as { count: number };
|
||||
if (pending.count === 0) {
|
||||
this.database.prepare("UPDATE asset_cleanup_requests SET status = 'completed' WHERE request_id = ? AND status = 'queued'")
|
||||
.run(request.request_id);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
this.database.prepare("DELETE FROM managed_files WHERE file_id = ?").run(row.managed_file_id);
|
||||
if (this.tableExists("local_backend_storage_state")) this.decrementManagedCapacity(row.byte_size);
|
||||
}
|
||||
this.database.prepare(`
|
||||
UPDATE file_cleanup_queue SET status = 'completed', completed_at = ?, last_error = NULL
|
||||
WHERE cleanup_id = ?
|
||||
`).run(iso(this.clock()), row.cleanup_id);
|
||||
});
|
||||
transaction.immediate();
|
||||
completed += 1;
|
||||
} catch {
|
||||
this.database.prepare(`
|
||||
UPDATE file_cleanup_queue SET status = 'failed', last_error = 'physical_file_cleanup_failed'
|
||||
WHERE cleanup_id = ?
|
||||
`).run(row.cleanup_id);
|
||||
failed += 1;
|
||||
}
|
||||
}
|
||||
return { completed, failed };
|
||||
}
|
||||
|
||||
private queueManagedFiles(projectId: string, now: number) {
|
||||
if (!["project_resource_files", "managed_files", "project_asset_refs", "file_cleanup_queue"]
|
||||
.every((table) => this.tableExists(table))) return;
|
||||
const files = this.database.prepare(`
|
||||
SELECT mf.file_id, mf.relative_path, mf.byte_size
|
||||
FROM project_resource_files prf
|
||||
JOIN managed_files mf ON mf.file_id = prf.managed_file_id
|
||||
WHERE prf.project_id = ? AND mf.status = 'committed'
|
||||
ORDER BY mf.file_id
|
||||
`).all(projectId) as Array<{ byte_size: number; file_id: string; relative_path: string }>;
|
||||
for (const file of files) {
|
||||
this.database.prepare("DELETE FROM project_asset_refs WHERE reference_id = ?")
|
||||
.run(`project:${projectId}:${file.file_id}`);
|
||||
const otherProject = this.database.prepare(`
|
||||
SELECT 1 FROM project_resource_files prf
|
||||
JOIN projects p ON p.project_id = prf.project_id
|
||||
WHERE prf.managed_file_id = ? AND prf.project_id <> ? AND p.status <> 'purged'
|
||||
LIMIT 1
|
||||
`).get(file.file_id, projectId);
|
||||
const otherReference = this.database.prepare("SELECT 1 FROM project_asset_refs WHERE managed_file_id = ? LIMIT 1").get(file.file_id);
|
||||
if (otherProject || otherReference) continue;
|
||||
this.database.prepare("UPDATE managed_files SET status = 'purged', purged_at = ? WHERE file_id = ? AND status = 'committed'")
|
||||
.run(iso(now), file.file_id);
|
||||
this.database.prepare(`
|
||||
INSERT OR IGNORE INTO file_cleanup_queue (
|
||||
cleanup_id, managed_file_id, relative_path, byte_size, counts_toward_managed,
|
||||
reason, status, created_at, completed_at, last_error
|
||||
) VALUES (?, ?, ?, ?, 1, 'purge', 'pending', ?, NULL, NULL)
|
||||
`).run(randomUUID(), file.file_id, file.relative_path, file.byte_size, iso(now));
|
||||
}
|
||||
}
|
||||
|
||||
private tableExists(name: string) {
|
||||
return Boolean(this.database.prepare("SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?").get(name));
|
||||
}
|
||||
|
||||
private resolveManagedPath(relativePath: string) {
|
||||
const path = resolve(this.dataRoot!, relativePath);
|
||||
const scoped = relative(this.dataRoot!, path);
|
||||
if (!scoped || scoped.startsWith("..") || isAbsolute(scoped)) throw new Error("managed_path_outside_data_root");
|
||||
return path;
|
||||
}
|
||||
|
||||
private decrementManagedCapacity(bytes: number) {
|
||||
const current = this.database.prepare(`
|
||||
SELECT managed_content_bytes FROM local_backend_storage_state WHERE singleton = 1
|
||||
`).get() as { managed_content_bytes: number } | undefined;
|
||||
if (!current) return;
|
||||
const managed = Math.max(0, current.managed_content_bytes - bytes);
|
||||
const notice = managed < 4_294_967_296 ? "normal" : managed < 4_831_838_208 ? "warning" : "critical";
|
||||
const reservations = this.tableExists("storage_reservations")
|
||||
? (this.database.prepare("SELECT COALESCE(SUM(projected_bytes), 0) AS bytes FROM storage_reservations WHERE status = 'active'").get() as { bytes: number }).bytes
|
||||
: 0;
|
||||
const status = managed + reservations >= 5_368_709_120 ? "full" : "active";
|
||||
this.database.prepare(`
|
||||
UPDATE local_backend_storage_state
|
||||
SET managed_content_bytes = ?, capacity_notice_level = ?, storage_status = ?, measured_at = ?, version = version + 1
|
||||
WHERE singleton = 1
|
||||
`).run(managed, notice, status, iso(this.clock()));
|
||||
}
|
||||
}
|
||||
@@ -4,6 +4,7 @@ import { join } from "node:path";
|
||||
import { WorkerAiCallGate } from "./ai-call-gate.js";
|
||||
import { readConfiguredLocalDataRoot } from "./runtime-config.js";
|
||||
import { RetentionCleanup } from "./retention-cleanup.js";
|
||||
import { ProjectPurgeCleanup } from "./project-purge-cleanup.js";
|
||||
import { StructuredJsonlLogger } from "./structured-log.js";
|
||||
import { WorkerStorageStatus } from "./storage-status.js";
|
||||
import { attachWorkerSupervisorControl, initializeWorkerCredentialClient, receiveWorkerCredentials } from "./supervisor-channel.js";
|
||||
@@ -28,11 +29,13 @@ if (!workerPort && process.argv.includes("--dada-credential-stdin")) {
|
||||
const keepAlive = setInterval(() => undefined, 30_000);
|
||||
let storage: WorkerStorageStatus | undefined;
|
||||
let retention: RetentionCleanup | undefined;
|
||||
let projectCleanup: ProjectPurgeCleanup | undefined;
|
||||
let retentionTimer: ReturnType<typeof setInterval> | undefined;
|
||||
const control = attachWorkerSupervisorControl(controlPipe, () => {
|
||||
clearInterval(keepAlive);
|
||||
if (retentionTimer) clearInterval(retentionTimer);
|
||||
retention?.close();
|
||||
projectCleanup?.close();
|
||||
storage?.close();
|
||||
});
|
||||
let storageStatus: "active" | "unavailable" = "active";
|
||||
@@ -41,9 +44,11 @@ if (!workerPort && process.argv.includes("--dada-credential-stdin")) {
|
||||
const databasePath = join(dataRoot, "db", "dada.sqlite3");
|
||||
storage = new WorkerStorageStatus(databasePath);
|
||||
retention = new RetentionCleanup({ databasePath });
|
||||
projectCleanup = new ProjectPurgeCleanup({ dataRoot, databasePath });
|
||||
const runRetentionCleanup = () => {
|
||||
try {
|
||||
retention?.purgeExpired();
|
||||
projectCleanup?.run();
|
||||
} catch {
|
||||
storageStatus = "unavailable";
|
||||
storage?.markLogUnavailable();
|
||||
|
||||
Reference in New Issue
Block a user