feat: classify sensitive columns locally
This commit is contained in:
+34
-8
@@ -3,6 +3,7 @@ import cors from "@fastify/cors";
|
||||
import cookie from "@fastify/cookie";
|
||||
import rateLimit from "@fastify/rate-limit";
|
||||
import { dirname, isAbsolute, join } from "node:path";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { tmpdir } from "node:os";
|
||||
import type { AppConfig } from "./config.js";
|
||||
import { ThtRunner } from "./tht/tht-runner.js";
|
||||
@@ -57,8 +58,11 @@ import { metadataGenerationModelRoutes } from "./routes/metadata-generation-mode
|
||||
import { catalogDescriptionConsolidationRoutes } from "./routes/catalog-description-consolidation.js";
|
||||
import { PythonModelCompleter, type ModelCompleter } from "./catalog/model-completer.js";
|
||||
import { DescriptionGenerationWorker } from "./catalog/description-generation-worker.js";
|
||||
import { SensitiveDataSuggester } from "./catalog/sensitive-data-suggester.js";
|
||||
import { SensitiveDataSuggestionRunner } from "./catalog/sensitive-data-suggestion-runner.js";
|
||||
import { SensitivityAnalysisService } from "./catalog/sensitivity-analysis-service.js";
|
||||
import { SensitivityAnalysisRunner } from "./catalog/sensitivity-analysis-runner.js";
|
||||
import { SensitivityClassifier, type LocalNerDetector, type SensitivityValueSource } from "./catalog/sensitivity-classifier.js";
|
||||
import { ConcreteSensitivityValueSource } from "./catalog/sensitivity-value-source.js";
|
||||
import { PythonLocalNerDetector } from "./catalog/local-ner-detector.js";
|
||||
import {
|
||||
ConcreteDescriptionSourceSampler,
|
||||
type DescriptionSourceSampler,
|
||||
@@ -93,6 +97,8 @@ export interface BuildAppDeps {
|
||||
runtimeModelCatalog?: RuntimeModelCatalog;
|
||||
modelCompleter?: ModelCompleter;
|
||||
descriptionSourceSampler?: DescriptionSourceSampler;
|
||||
sensitivityValueSource?: SensitivityValueSource;
|
||||
localNerDetector?: LocalNerDetector;
|
||||
workspaceRuntimeSupport?: (workspace: WorkspaceDescriptor) => boolean;
|
||||
maintenanceBarrier?: MaintenanceBarrier;
|
||||
piManagement?: PiManagementService;
|
||||
@@ -194,12 +200,24 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
|
||||
catalogOperationCoordinator,
|
||||
descriptionSourceSampler,
|
||||
);
|
||||
const sensitiveDataSuggester = new SensitiveDataSuggester(
|
||||
const sensitivityValueSource = deps?.sensitivityValueSource
|
||||
?? new ConcreteSensitivityValueSource(catalogPostgresAccess, workspaceSecretStore);
|
||||
const configuredNerWorker = config.sensitivityNer?.workerScript
|
||||
?? fileURLToPath(new URL("../python/sensitivity_ner_worker.py", import.meta.url));
|
||||
const localNerDetector = deps?.localNerDetector ?? (config.sensitivityNer
|
||||
? new PythonLocalNerDetector({
|
||||
pythonExecutable: config.sensitivityNer.pythonExecutable,
|
||||
workerScript: configuredNerWorker,
|
||||
modelPath: config.sensitivityNer.modelPath,
|
||||
cwd: dirname(configuredNerWorker),
|
||||
threads: config.sensitivityNer.threads,
|
||||
})
|
||||
: undefined);
|
||||
const sensitiveDataSuggester = new SensitivityAnalysisService(
|
||||
catalogRepository,
|
||||
metadataGenerationModels,
|
||||
modelCompleter,
|
||||
new SensitivityClassifier(sensitivityValueSource, localNerDetector),
|
||||
);
|
||||
const sensitiveDataSuggestionRunner = new SensitiveDataSuggestionRunner(
|
||||
const sensitivityAnalysisRunner = new SensitivityAnalysisRunner(
|
||||
catalogRepository,
|
||||
sensitiveDataSuggester,
|
||||
);
|
||||
@@ -235,12 +253,20 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
|
||||
);
|
||||
app.addHook("onReady", async () => { await catalogSyncWorker.initialize(); });
|
||||
app.addHook("onReady", async () => { await descriptionGenerationWorker.initialize(); });
|
||||
app.addHook("onReady", async () => { await sensitiveDataSuggestionRunner.initialize(); });
|
||||
app.addHook("onReady", async () => { await sensitivityAnalysisRunner.initialize(); });
|
||||
if (localNerDetector?.warmup) {
|
||||
app.addHook("onReady", async () => {
|
||||
void localNerDetector.warmup?.().catch(() => undefined);
|
||||
});
|
||||
}
|
||||
if (!deps?.catalogRepository && catalogRepository.close) {
|
||||
app.addHook("onClose", async () => { await catalogRepository.close?.(); });
|
||||
}
|
||||
app.addHook("onClose", async () => { await catalogSyncWorker.stop(); });
|
||||
app.addHook("onClose", async () => { await descriptionGenerationWorker.stop(); });
|
||||
if (localNerDetector?.close) {
|
||||
app.addHook("onClose", async () => { await localNerDetector.close?.(); });
|
||||
}
|
||||
const workspaceDiagnoser = deps?.workspaceDiagnoser
|
||||
?? createProductionWorkspaceDiagnoser(config.workspaceDiagnosticTimeoutMs, undefined, {
|
||||
internalQdrantUrl: config.internalQdrantUrl,
|
||||
@@ -483,7 +509,7 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
|
||||
catalogDescriptionGenerationRoutes(app, {
|
||||
repository: catalogRepository,
|
||||
worker: descriptionGenerationWorker,
|
||||
sensitiveDataSuggestionRunner,
|
||||
sensitivityAnalysisRunner,
|
||||
});
|
||||
settingsRoutes(app, { cfg: config, getSettings });
|
||||
piManagementRoutes(app, { service: piManagement });
|
||||
|
||||
@@ -0,0 +1,254 @@
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { spawn, type ChildProcessWithoutNullStreams } from "node:child_process";
|
||||
import { tmpdir } from "node:os";
|
||||
import { z } from "zod";
|
||||
import type {
|
||||
LocalNerCandidate,
|
||||
LocalNerDetector,
|
||||
LocalNerEvidence,
|
||||
} from "./sensitivity-classifier.js";
|
||||
|
||||
const MAX_LINE_BYTES = 64 * 1024;
|
||||
const candidateSchema = z.object({
|
||||
columnId: z.uuid(),
|
||||
text: z.string().min(1).max(500),
|
||||
}).strict();
|
||||
const workerMessageSchema = z.union([
|
||||
z.object({ ready: z.literal(true) }).strict(),
|
||||
z.object({
|
||||
id: z.uuid(),
|
||||
ok: z.literal(true),
|
||||
evidence: z.array(z.object({
|
||||
columnId: z.uuid(),
|
||||
label: z.string().min(1).max(80),
|
||||
confidence: z.number().min(0).max(1),
|
||||
}).strict()).max(1_000),
|
||||
}).strict(),
|
||||
z.object({ id: z.uuid(), ok: z.literal(false), error: z.string().min(1).max(80) }).strict(),
|
||||
]);
|
||||
|
||||
export class LocalNerUnavailableError extends Error {
|
||||
constructor() {
|
||||
super("local NER is unavailable");
|
||||
this.name = "LocalNerUnavailableError";
|
||||
}
|
||||
}
|
||||
|
||||
interface PendingRequest {
|
||||
resolve: (value: readonly LocalNerEvidence[]) => void;
|
||||
reject: (error: Error) => void;
|
||||
timer: ReturnType<typeof setTimeout>;
|
||||
signal: AbortSignal;
|
||||
cancel: () => void;
|
||||
}
|
||||
|
||||
/** Persistent JSONL adapter for the optional, CPU-only Python NER worker. */
|
||||
export class PythonLocalNerDetector implements LocalNerDetector {
|
||||
private child?: ChildProcessWithoutNullStreams;
|
||||
private ready?: Promise<void>;
|
||||
private readyResolve?: () => void;
|
||||
private readyReject?: (error: Error) => void;
|
||||
private workerReady = false;
|
||||
private stdout = "";
|
||||
private readonly pending = new Map<string, PendingRequest>();
|
||||
|
||||
constructor(private readonly options: {
|
||||
pythonExecutable: string;
|
||||
workerScript: string;
|
||||
modelPath: string;
|
||||
cwd: string;
|
||||
threads?: number;
|
||||
startupTimeoutMs?: number;
|
||||
}) {}
|
||||
|
||||
async warmup(): Promise<void> {
|
||||
await this.ensureStarted();
|
||||
}
|
||||
|
||||
isReady(): boolean {
|
||||
return this.workerReady
|
||||
&& this.child !== undefined
|
||||
&& this.child.exitCode === null
|
||||
&& this.child.signalCode === null;
|
||||
}
|
||||
|
||||
async detect(
|
||||
candidates: readonly LocalNerCandidate[],
|
||||
signal: AbortSignal,
|
||||
deadline: number,
|
||||
): Promise<readonly LocalNerEvidence[]> {
|
||||
const parsed = z.array(candidateSchema).min(1).max(128).parse(candidates);
|
||||
if (signal.aborted || deadline <= Date.now()) throw new LocalNerUnavailableError();
|
||||
await this.ensureStartedWithin(signal, deadline);
|
||||
if (!this.child || this.child.exitCode !== null || this.child.signalCode !== null) {
|
||||
throw new LocalNerUnavailableError();
|
||||
}
|
||||
const id = randomUUID();
|
||||
return await new Promise<readonly LocalNerEvidence[]>((resolve, reject) => {
|
||||
const fail = () => {
|
||||
this.finishPending(id);
|
||||
reject(new LocalNerUnavailableError());
|
||||
this.stopWorker();
|
||||
};
|
||||
const timer = setTimeout(fail, Math.max(1, Math.floor(deadline - Date.now())));
|
||||
const cancel = fail;
|
||||
const pending: PendingRequest = { resolve, reject, timer, signal, cancel };
|
||||
this.pending.set(id, pending);
|
||||
signal.addEventListener("abort", cancel, { once: true });
|
||||
this.child!.stdin.write(`${JSON.stringify({ id, candidates: parsed })}\n`, (error) => {
|
||||
if (error) fail();
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
async close(): Promise<void> {
|
||||
const child = this.child;
|
||||
if (!child || child.exitCode !== null || child.signalCode !== null) return;
|
||||
await new Promise<void>((resolve) => {
|
||||
child.once("close", () => resolve());
|
||||
child.kill("SIGTERM");
|
||||
setTimeout(() => {
|
||||
if (child.exitCode === null && child.signalCode === null) child.kill("SIGKILL");
|
||||
}, 250).unref();
|
||||
});
|
||||
}
|
||||
|
||||
private async ensureStarted(): Promise<void> {
|
||||
if (this.ready) return await this.ready;
|
||||
this.ready = new Promise<void>((resolve, reject) => {
|
||||
this.readyResolve = resolve;
|
||||
this.readyReject = reject;
|
||||
});
|
||||
const threads = String(this.options.threads ?? 2);
|
||||
const inheritedRuntimeEnvironment = Object.fromEntries([
|
||||
"PATH", "SystemRoot", "WINDIR", "PATHEXT", "TMPDIR", "TEMP", "TMP", "LANG", "LC_ALL",
|
||||
].flatMap((name) => process.env[name] === undefined ? [] : [[name, process.env[name]!]]));
|
||||
const child = spawn(this.options.pythonExecutable, [
|
||||
"-I",
|
||||
"-B",
|
||||
this.options.workerScript,
|
||||
"--model",
|
||||
this.options.modelPath,
|
||||
"--threads",
|
||||
threads,
|
||||
], {
|
||||
cwd: this.options.cwd,
|
||||
stdio: ["pipe", "pipe", "pipe"],
|
||||
env: {
|
||||
...inheritedRuntimeEnvironment,
|
||||
HOME: process.env.HOME ?? tmpdir(),
|
||||
CUDA_VISIBLE_DEVICES: "",
|
||||
HIP_VISIBLE_DEVICES: "",
|
||||
HF_HUB_OFFLINE: "1",
|
||||
HF_HUB_DISABLE_TELEMETRY: "1",
|
||||
TRANSFORMERS_OFFLINE: "1",
|
||||
TOKENIZERS_PARALLELISM: "false",
|
||||
PYTHONNOUSERSITE: "1",
|
||||
OMP_NUM_THREADS: threads,
|
||||
MKL_NUM_THREADS: threads,
|
||||
OPENBLAS_NUM_THREADS: threads,
|
||||
HTTP_PROXY: "",
|
||||
HTTPS_PROXY: "",
|
||||
ALL_PROXY: "",
|
||||
NO_PROXY: "*",
|
||||
},
|
||||
});
|
||||
this.child = child;
|
||||
child.stdout.setEncoding("utf8");
|
||||
child.stdout.on("data", (chunk: string) => this.receive(chunk));
|
||||
child.stderr.resume();
|
||||
child.once("error", () => this.failWorker());
|
||||
child.once("close", () => this.failWorker());
|
||||
const startupTimer = setTimeout(() => this.failWorker(), this.options.startupTimeoutMs ?? 120_000);
|
||||
startupTimer.unref();
|
||||
try {
|
||||
await this.ready;
|
||||
} finally {
|
||||
clearTimeout(startupTimer);
|
||||
}
|
||||
}
|
||||
|
||||
private async ensureStartedWithin(signal: AbortSignal, deadline: number): Promise<void> {
|
||||
const started = this.ensureStarted();
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
let settled = false;
|
||||
const finish = (error?: Error, stopWorker = false) => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
clearTimeout(timer);
|
||||
signal.removeEventListener("abort", cancel);
|
||||
if (stopWorker) this.failWorker();
|
||||
if (error) reject(error);
|
||||
else resolve();
|
||||
};
|
||||
const cancel = () => finish(new LocalNerUnavailableError(), true);
|
||||
const timer = setTimeout(cancel, Math.max(1, Math.floor(deadline - Date.now())));
|
||||
signal.addEventListener("abort", cancel, { once: true });
|
||||
void started.then(
|
||||
() => finish(),
|
||||
() => finish(new LocalNerUnavailableError()),
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
private receive(chunk: string): void {
|
||||
this.stdout += chunk;
|
||||
if (Buffer.byteLength(this.stdout, "utf8") > MAX_LINE_BYTES) {
|
||||
this.failWorker();
|
||||
return;
|
||||
}
|
||||
let newline: number;
|
||||
while ((newline = this.stdout.indexOf("\n")) >= 0) {
|
||||
const line = this.stdout.slice(0, newline);
|
||||
this.stdout = this.stdout.slice(newline + 1);
|
||||
if (!line) continue;
|
||||
try {
|
||||
const message = workerMessageSchema.parse(JSON.parse(line));
|
||||
if ("ready" in message) {
|
||||
this.workerReady = true;
|
||||
this.readyResolve?.();
|
||||
this.readyResolve = undefined;
|
||||
this.readyReject = undefined;
|
||||
continue;
|
||||
}
|
||||
const pending = this.pending.get(message.id);
|
||||
if (!pending) continue;
|
||||
this.finishPending(message.id);
|
||||
if (message.ok) pending.resolve(message.evidence);
|
||||
else pending.reject(new LocalNerUnavailableError());
|
||||
} catch {
|
||||
this.failWorker();
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private finishPending(id: string): void {
|
||||
const pending = this.pending.get(id);
|
||||
if (!pending) return;
|
||||
clearTimeout(pending.timer);
|
||||
pending.signal.removeEventListener("abort", pending.cancel);
|
||||
this.pending.delete(id);
|
||||
}
|
||||
|
||||
private stopWorker(): void {
|
||||
const child = this.child;
|
||||
if (child && child.exitCode === null && child.signalCode === null) child.kill("SIGTERM");
|
||||
}
|
||||
|
||||
private failWorker(): void {
|
||||
const error = new LocalNerUnavailableError();
|
||||
this.readyReject?.(error);
|
||||
this.readyResolve = undefined;
|
||||
this.readyReject = undefined;
|
||||
for (const [id, pending] of this.pending) {
|
||||
this.finishPending(id);
|
||||
pending.reject(error);
|
||||
}
|
||||
this.stopWorker();
|
||||
this.child = undefined;
|
||||
this.ready = undefined;
|
||||
this.workerReady = false;
|
||||
this.stdout = "";
|
||||
}
|
||||
}
|
||||
@@ -35,10 +35,10 @@ import {
|
||||
type DescriptionGenerationRun,
|
||||
type DescriptionGenerationRunUpdate,
|
||||
type DescriptionGenerationScope,
|
||||
type SensitiveDataSuggestionEvent,
|
||||
type SensitiveDataSuggestionRun,
|
||||
type SensitiveDataSuggestionRunUpdate,
|
||||
type SensitiveDataSuggestionScope,
|
||||
type SensitivityAnalysisEvent,
|
||||
type SensitivityAnalysisRun,
|
||||
type SensitivityAnalysisRunUpdate,
|
||||
type SensitivityAnalysisScope,
|
||||
type TableSyncRepositoryResult,
|
||||
type WorkspaceDatabase,
|
||||
} from "./types.js";
|
||||
@@ -56,8 +56,8 @@ export class MemoryCatalogRepository implements CatalogRepository {
|
||||
private readonly logicalRelationships = new Map<string, CatalogLogicalRelationship>();
|
||||
private readonly descriptionGenerationRuns = new Map<string, DescriptionGenerationRun>();
|
||||
private readonly descriptionGenerationEvents = new Map<string, DescriptionGenerationEvent[]>();
|
||||
private readonly sensitiveDataSuggestionRuns = new Map<string, SensitiveDataSuggestionRun>();
|
||||
private readonly sensitiveDataSuggestionEvents = new Map<string, SensitiveDataSuggestionEvent[]>();
|
||||
private readonly sensitivityAnalysisRuns = new Map<string, SensitivityAnalysisRun>();
|
||||
private readonly sensitivityAnalysisEvents = new Map<string, SensitivityAnalysisEvent[]>();
|
||||
private readonly syncRuns = new Map<string, CatalogSyncRun>();
|
||||
private readonly syncEvents = new Map<string, CatalogSyncEvent[]>();
|
||||
|
||||
@@ -183,10 +183,10 @@ export class MemoryCatalogRepository implements CatalogRepository {
|
||||
this.descriptionGenerationRuns.delete(runId);
|
||||
this.descriptionGenerationEvents.delete(runId);
|
||||
}
|
||||
for (const [runId, run] of this.sensitiveDataSuggestionRuns) {
|
||||
for (const [runId, run] of this.sensitivityAnalysisRuns) {
|
||||
if (run.databaseId !== id) continue;
|
||||
this.sensitiveDataSuggestionRuns.delete(runId);
|
||||
this.sensitiveDataSuggestionEvents.delete(runId);
|
||||
this.sensitivityAnalysisRuns.delete(runId);
|
||||
this.sensitivityAnalysisEvents.delete(runId);
|
||||
}
|
||||
return this.records.delete(id);
|
||||
}
|
||||
@@ -433,93 +433,96 @@ export class MemoryCatalogRepository implements CatalogRepository {
|
||||
.map((event) => structuredClone(event));
|
||||
}
|
||||
|
||||
async createSensitiveDataSuggestionRun(
|
||||
async createSensitivityAnalysisRun(
|
||||
databaseId: string,
|
||||
scope: SensitiveDataSuggestionScope,
|
||||
modelId: string,
|
||||
): Promise<SensitiveDataSuggestionRun> {
|
||||
scope: SensitivityAnalysisScope,
|
||||
origin: { engine: "llm"; modelId: string } | { engine: "local"; policyVersion: string },
|
||||
): Promise<SensitivityAnalysisRun> {
|
||||
const now = new Date().toISOString();
|
||||
const run: SensitiveDataSuggestionRun = {
|
||||
const run: SensitivityAnalysisRun = {
|
||||
id: randomUUID(),
|
||||
databaseId,
|
||||
scope,
|
||||
modelId,
|
||||
engine: origin.engine,
|
||||
modelId: origin.engine === "llm" ? origin.modelId : null,
|
||||
policyVersion: origin.engine === "local" ? origin.policyVersion : null,
|
||||
status: "running",
|
||||
total: 0,
|
||||
suggestedSensitive: 0,
|
||||
suggestedNonSensitive: 0,
|
||||
inputTokens: 0,
|
||||
cacheReadTokens: 0,
|
||||
outputTokens: 0,
|
||||
suggestedNonSensitive: 0,
|
||||
unknown: 0,
|
||||
inputTokens: 0,
|
||||
cacheReadTokens: 0,
|
||||
outputTokens: 0,
|
||||
createdAt: now,
|
||||
startedAt: now,
|
||||
updatedAt: now,
|
||||
finishedAt: null,
|
||||
errorSummary: null,
|
||||
};
|
||||
this.sensitiveDataSuggestionRuns.set(run.id, run);
|
||||
this.sensitivityAnalysisRuns.set(run.id, run);
|
||||
return structuredClone(run);
|
||||
}
|
||||
|
||||
async getSensitiveDataSuggestionRun(
|
||||
async getSensitivityAnalysisRun(
|
||||
runId: string,
|
||||
): Promise<SensitiveDataSuggestionRun | undefined> {
|
||||
const run = this.sensitiveDataSuggestionRuns.get(runId);
|
||||
): Promise<SensitivityAnalysisRun | undefined> {
|
||||
const run = this.sensitivityAnalysisRuns.get(runId);
|
||||
return run ? structuredClone(run) : undefined;
|
||||
}
|
||||
|
||||
async listSensitiveDataSuggestionRuns(limit = 50): Promise<SensitiveDataSuggestionRun[]> {
|
||||
return [...this.sensitiveDataSuggestionRuns.values()]
|
||||
async listSensitivityAnalysisRuns(limit = 50): Promise<SensitivityAnalysisRun[]> {
|
||||
return [...this.sensitivityAnalysisRuns.values()]
|
||||
.sort((a, b) => b.createdAt.localeCompare(a.createdAt) || b.id.localeCompare(a.id))
|
||||
.slice(0, limit)
|
||||
.map((run) => structuredClone(run));
|
||||
}
|
||||
|
||||
async interruptActiveSensitiveDataSuggestionRuns(
|
||||
async interruptActiveSensitivityAnalysisRuns(
|
||||
errorSummary: string,
|
||||
): Promise<SensitiveDataSuggestionRun[]> {
|
||||
const interrupted: SensitiveDataSuggestionRun[] = [];
|
||||
for (const run of this.sensitiveDataSuggestionRuns.values()) {
|
||||
): Promise<SensitivityAnalysisRun[]> {
|
||||
const interrupted: SensitivityAnalysisRun[] = [];
|
||||
for (const run of this.sensitivityAnalysisRuns.values()) {
|
||||
if (run.status !== "running") continue;
|
||||
const now = new Date().toISOString();
|
||||
const updated: SensitiveDataSuggestionRun = {
|
||||
const updated: SensitivityAnalysisRun = {
|
||||
...run,
|
||||
status: "interrupted",
|
||||
updatedAt: now,
|
||||
finishedAt: now,
|
||||
errorSummary,
|
||||
};
|
||||
this.sensitiveDataSuggestionRuns.set(run.id, updated);
|
||||
this.sensitivityAnalysisRuns.set(run.id, updated);
|
||||
interrupted.push(structuredClone(updated));
|
||||
}
|
||||
return interrupted;
|
||||
}
|
||||
|
||||
async updateSensitiveDataSuggestionRun(
|
||||
async updateSensitivityAnalysisRun(
|
||||
runId: string,
|
||||
update: SensitiveDataSuggestionRunUpdate,
|
||||
): Promise<SensitiveDataSuggestionRun | undefined> {
|
||||
const current = this.sensitiveDataSuggestionRuns.get(runId);
|
||||
update: SensitivityAnalysisRunUpdate,
|
||||
): Promise<SensitivityAnalysisRun | undefined> {
|
||||
const current = this.sensitivityAnalysisRuns.get(runId);
|
||||
if (!current) return undefined;
|
||||
const updated = {
|
||||
...current,
|
||||
...structuredClone(update),
|
||||
updatedAt: new Date().toISOString(),
|
||||
};
|
||||
this.sensitiveDataSuggestionRuns.set(runId, updated);
|
||||
this.sensitivityAnalysisRuns.set(runId, updated);
|
||||
return structuredClone(updated);
|
||||
}
|
||||
|
||||
async appendSensitiveDataSuggestionEvent(
|
||||
async appendSensitivityAnalysisEvent(
|
||||
runId: string,
|
||||
level: SensitiveDataSuggestionEvent["level"],
|
||||
level: SensitivityAnalysisEvent["level"],
|
||||
message: string,
|
||||
): Promise<SensitiveDataSuggestionEvent> {
|
||||
if (!this.sensitiveDataSuggestionRuns.has(runId)) {
|
||||
throw new CatalogConflictError("Sensitive Data Suggestion Run does not exist");
|
||||
): Promise<SensitivityAnalysisEvent> {
|
||||
if (!this.sensitivityAnalysisRuns.has(runId)) {
|
||||
throw new CatalogConflictError("Sensitivity Analysis Run does not exist");
|
||||
}
|
||||
const events = this.sensitiveDataSuggestionEvents.get(runId) ?? [];
|
||||
const event: SensitiveDataSuggestionEvent = {
|
||||
const events = this.sensitivityAnalysisEvents.get(runId) ?? [];
|
||||
const event: SensitivityAnalysisEvent = {
|
||||
runId,
|
||||
sequence: events.length + 1,
|
||||
level,
|
||||
@@ -527,15 +530,15 @@ export class MemoryCatalogRepository implements CatalogRepository {
|
||||
createdAt: new Date().toISOString(),
|
||||
};
|
||||
events.push(event);
|
||||
this.sensitiveDataSuggestionEvents.set(runId, events);
|
||||
this.sensitivityAnalysisEvents.set(runId, events);
|
||||
return structuredClone(event);
|
||||
}
|
||||
|
||||
async listSensitiveDataSuggestionEvents(
|
||||
async listSensitivityAnalysisEvents(
|
||||
runId: string,
|
||||
afterSequence = 0,
|
||||
): Promise<SensitiveDataSuggestionEvent[]> {
|
||||
return (this.sensitiveDataSuggestionEvents.get(runId) ?? [])
|
||||
): Promise<SensitivityAnalysisEvent[]> {
|
||||
return (this.sensitivityAnalysisEvents.get(runId) ?? [])
|
||||
.filter((event) => event.sequence > afterSequence)
|
||||
.map((event) => structuredClone(event));
|
||||
}
|
||||
|
||||
@@ -9,10 +9,11 @@ import * as catalogSchemaSyncMigration from "./migrations/003_catalog_schema_syn
|
||||
import * as catalogRuntimeSequencePrivilegesMigration from "./migrations/004_catalog_runtime_sequence_privileges.js";
|
||||
import * as descriptionGenerationRunsMigration from "./migrations/005_description_generation_runs.js";
|
||||
import * as sensitiveDataFlagMigration from "./migrations/006_sensitive_data_flag.js";
|
||||
import * as sensitiveDataSuggestionRunsMigration from "./migrations/007_sensitive_data_suggestion_runs.js";
|
||||
import * as sensitivityAnalysisRunsMigration from "./migrations/007_sensitive_data_suggestion_runs.js";
|
||||
import * as catalogLogicalRelationshipsMigration from "./migrations/008_catalog_logical_relationships.js";
|
||||
import * as aiTokenUsageMigration from "./migrations/009_ai_token_usage.js";
|
||||
import * as canonicalModelIdsMigration from "./migrations/010_canonical_model_ids.js";
|
||||
import * as localSensitivityAnalysisMigration from "./migrations/011_local_sensitivity_analysis.js";
|
||||
|
||||
const connectionString = process.env.THT_CATALOG_MIGRATOR_DATABASE_URL;
|
||||
const host = process.env.THT_CATALOG_DB_HOST;
|
||||
@@ -44,10 +45,11 @@ const provider: MigrationProvider = {
|
||||
"004_catalog_runtime_sequence_privileges": catalogRuntimeSequencePrivilegesMigration,
|
||||
"005_description_generation_runs": descriptionGenerationRunsMigration,
|
||||
"006_sensitive_data_flag": sensitiveDataFlagMigration,
|
||||
"007_sensitive_data_suggestion_runs": sensitiveDataSuggestionRunsMigration,
|
||||
"007_sensitive_data_suggestion_runs": sensitivityAnalysisRunsMigration,
|
||||
"008_catalog_logical_relationships": catalogLogicalRelationshipsMigration,
|
||||
"009_ai_token_usage": aiTokenUsageMigration,
|
||||
"010_canonical_model_ids": canonicalModelIdsMigration,
|
||||
"011_local_sensitivity_analysis": localSensitivityAnalysisMigration,
|
||||
};
|
||||
},
|
||||
};
|
||||
|
||||
@@ -0,0 +1,42 @@
|
||||
import { sql, type Kysely } from "kysely";
|
||||
import type { CatalogDatabase } from "../repository.js";
|
||||
|
||||
export async function up(db: Kysely<CatalogDatabase>): Promise<void> {
|
||||
await sql.raw(`alter table sensitive_data_suggestion_runs
|
||||
alter column model_id drop not null,
|
||||
add column engine text not null default 'llm',
|
||||
add column policy_version text,
|
||||
add column unknown integer not null default 0,
|
||||
drop constraint sensitive_data_suggestion_runs_counters_check,
|
||||
add constraint sensitive_data_suggestion_runs_counters_check
|
||||
check (total >= 0
|
||||
and suggested_sensitive >= 0
|
||||
and suggested_non_sensitive >= 0
|
||||
and unknown >= 0
|
||||
and suggested_sensitive + suggested_non_sensitive + unknown <= total),
|
||||
add constraint sensitive_data_suggestion_runs_engine_check
|
||||
check (engine in ('llm', 'local')),
|
||||
add constraint sensitive_data_suggestion_runs_origin_check
|
||||
check ((engine = 'llm' and model_id is not null and policy_version is null)
|
||||
or (engine = 'local' and model_id is null
|
||||
and policy_version ~ '^[a-z][a-z0-9._-]{0,63}$'))`).execute(db);
|
||||
}
|
||||
|
||||
export async function down(db: Kysely<CatalogDatabase>): Promise<void> {
|
||||
await sql.raw(`alter table sensitive_data_suggestion_runs
|
||||
drop constraint sensitive_data_suggestion_runs_origin_check,
|
||||
drop constraint sensitive_data_suggestion_runs_engine_check,
|
||||
drop constraint sensitive_data_suggestion_runs_counters_check`).execute(db);
|
||||
await sql.raw(`update sensitive_data_suggestion_runs
|
||||
set model_id = coalesce(model_id, 'local/sensitivity-v1')`).execute(db);
|
||||
await sql.raw(`alter table sensitive_data_suggestion_runs
|
||||
drop column unknown,
|
||||
drop column policy_version,
|
||||
drop column engine,
|
||||
alter column model_id set not null,
|
||||
add constraint sensitive_data_suggestion_runs_counters_check
|
||||
check (total >= 0
|
||||
and suggested_sensitive >= 0
|
||||
and suggested_non_sensitive >= 0
|
||||
and suggested_sensitive + suggested_non_sensitive <= total)`).execute(db);
|
||||
}
|
||||
@@ -45,10 +45,10 @@ import {
|
||||
type DescriptionGenerationScope,
|
||||
type ObservedCatalogTable,
|
||||
type ObservedSchemaSnapshot,
|
||||
type SensitiveDataSuggestionEvent,
|
||||
type SensitiveDataSuggestionRun,
|
||||
type SensitiveDataSuggestionRunUpdate,
|
||||
type SensitiveDataSuggestionScope,
|
||||
type SensitivityAnalysisEvent,
|
||||
type SensitivityAnalysisRun,
|
||||
type SensitivityAnalysisRunUpdate,
|
||||
type SensitivityAnalysisScope,
|
||||
type TableSyncRepositoryResult,
|
||||
type WorkspaceDatabase,
|
||||
} from "./types.js";
|
||||
@@ -189,15 +189,18 @@ interface DescriptionGenerationEventTable {
|
||||
createdAt: Timestamp;
|
||||
}
|
||||
|
||||
interface SensitiveDataSuggestionRunTable {
|
||||
interface SensitivityAnalysisRunTable {
|
||||
id: string;
|
||||
databaseId: string;
|
||||
scope: SensitiveDataSuggestionScope;
|
||||
modelId: string;
|
||||
status: SensitiveDataSuggestionRun["status"];
|
||||
scope: SensitivityAnalysisScope;
|
||||
engine: SensitivityAnalysisRun["engine"];
|
||||
modelId: string | null;
|
||||
policyVersion: string | null;
|
||||
status: SensitivityAnalysisRun["status"];
|
||||
total: number;
|
||||
suggestedSensitive: number;
|
||||
suggestedNonSensitive: number;
|
||||
unknown: number;
|
||||
inputTokens: number;
|
||||
cacheReadTokens: number;
|
||||
outputTokens: number;
|
||||
@@ -208,10 +211,10 @@ interface SensitiveDataSuggestionRunTable {
|
||||
errorSummary: string | null;
|
||||
}
|
||||
|
||||
interface SensitiveDataSuggestionEventTable {
|
||||
interface SensitivityAnalysisEventTable {
|
||||
runId: string;
|
||||
sequence: number;
|
||||
level: SensitiveDataSuggestionEvent["level"];
|
||||
level: SensitivityAnalysisEvent["level"];
|
||||
message: string;
|
||||
createdAt: Timestamp;
|
||||
}
|
||||
@@ -264,8 +267,9 @@ export interface CatalogDatabase {
|
||||
catalogLogicalRelationships: CatalogLogicalRelationshipTable;
|
||||
descriptionGenerationRuns: DescriptionGenerationRunTable;
|
||||
descriptionGenerationEvents: DescriptionGenerationEventTable;
|
||||
sensitiveDataSuggestionRuns: SensitiveDataSuggestionRunTable;
|
||||
sensitiveDataSuggestionEvents: SensitiveDataSuggestionEventTable;
|
||||
// Legacy physical table names retained for migration and storage compatibility.
|
||||
sensitiveDataSuggestionRuns: SensitivityAnalysisRunTable;
|
||||
sensitiveDataSuggestionEvents: SensitivityAnalysisEventTable;
|
||||
catalogSyncRuns: CatalogSyncRunTable;
|
||||
catalogSyncEvents: CatalogSyncEventTable;
|
||||
}
|
||||
@@ -402,9 +406,9 @@ function serializeDescriptionGenerationEvent(
|
||||
return { ...row, createdAt: new Date(row.createdAt).toISOString() };
|
||||
}
|
||||
|
||||
function serializeSensitiveDataSuggestionRun(
|
||||
row: Selectable<SensitiveDataSuggestionRunTable>,
|
||||
): SensitiveDataSuggestionRun {
|
||||
function serializeSensitivityAnalysisRun(
|
||||
row: Selectable<SensitivityAnalysisRunTable>,
|
||||
): SensitivityAnalysisRun {
|
||||
const stamp = (value: Date | string | null) => value === null ? null : new Date(value).toISOString();
|
||||
return {
|
||||
...row,
|
||||
@@ -415,9 +419,9 @@ function serializeSensitiveDataSuggestionRun(
|
||||
};
|
||||
}
|
||||
|
||||
function serializeSensitiveDataSuggestionEvent(
|
||||
row: Selectable<SensitiveDataSuggestionEventTable>,
|
||||
): SensitiveDataSuggestionEvent {
|
||||
function serializeSensitivityAnalysisEvent(
|
||||
row: Selectable<SensitivityAnalysisEventTable>,
|
||||
): SensitivityAnalysisEvent {
|
||||
return { ...row, createdAt: new Date(row.createdAt).toISOString() };
|
||||
}
|
||||
|
||||
@@ -938,52 +942,55 @@ export class KyselyCatalogRepository implements CatalogRepository {
|
||||
return rows.map(serializeDescriptionGenerationEvent);
|
||||
}
|
||||
|
||||
async createSensitiveDataSuggestionRun(
|
||||
async createSensitivityAnalysisRun(
|
||||
databaseId: string,
|
||||
scope: SensitiveDataSuggestionScope,
|
||||
modelId: string,
|
||||
): Promise<SensitiveDataSuggestionRun> {
|
||||
scope: SensitivityAnalysisScope,
|
||||
origin: { engine: "llm"; modelId: string } | { engine: "local"; policyVersion: string },
|
||||
): Promise<SensitivityAnalysisRun> {
|
||||
const row = await this.db.insertInto("sensitiveDataSuggestionRuns").values({
|
||||
id: randomUUID(),
|
||||
databaseId,
|
||||
scope,
|
||||
modelId,
|
||||
engine: origin.engine,
|
||||
modelId: origin.engine === "llm" ? origin.modelId : null,
|
||||
policyVersion: origin.engine === "local" ? origin.policyVersion : null,
|
||||
status: "running",
|
||||
total: 0,
|
||||
suggestedSensitive: 0,
|
||||
suggestedNonSensitive: 0,
|
||||
unknown: 0,
|
||||
inputTokens: 0,
|
||||
cacheReadTokens: 0,
|
||||
outputTokens: 0,
|
||||
finishedAt: null,
|
||||
errorSummary: null,
|
||||
}).returningAll().executeTakeFirstOrThrow();
|
||||
return serializeSensitiveDataSuggestionRun(row);
|
||||
return serializeSensitivityAnalysisRun(row);
|
||||
}
|
||||
|
||||
async getSensitiveDataSuggestionRun(
|
||||
async getSensitivityAnalysisRun(
|
||||
runId: string,
|
||||
): Promise<SensitiveDataSuggestionRun | undefined> {
|
||||
): Promise<SensitivityAnalysisRun | undefined> {
|
||||
const row = await this.db.selectFrom("sensitiveDataSuggestionRuns")
|
||||
.selectAll()
|
||||
.where("id", "=", runId)
|
||||
.executeTakeFirst();
|
||||
return row ? serializeSensitiveDataSuggestionRun(row) : undefined;
|
||||
return row ? serializeSensitivityAnalysisRun(row) : undefined;
|
||||
}
|
||||
|
||||
async listSensitiveDataSuggestionRuns(limit = 50): Promise<SensitiveDataSuggestionRun[]> {
|
||||
async listSensitivityAnalysisRuns(limit = 50): Promise<SensitivityAnalysisRun[]> {
|
||||
const rows = await this.db.selectFrom("sensitiveDataSuggestionRuns")
|
||||
.selectAll()
|
||||
.orderBy("createdAt", "desc")
|
||||
.orderBy("id", "desc")
|
||||
.limit(limit)
|
||||
.execute();
|
||||
return rows.map(serializeSensitiveDataSuggestionRun);
|
||||
return rows.map(serializeSensitivityAnalysisRun);
|
||||
}
|
||||
|
||||
async interruptActiveSensitiveDataSuggestionRuns(
|
||||
async interruptActiveSensitivityAnalysisRuns(
|
||||
errorSummary: string,
|
||||
): Promise<SensitiveDataSuggestionRun[]> {
|
||||
): Promise<SensitivityAnalysisRun[]> {
|
||||
const rows = await this.db.updateTable("sensitiveDataSuggestionRuns")
|
||||
.set({
|
||||
status: "interrupted",
|
||||
@@ -994,34 +1001,34 @@ export class KyselyCatalogRepository implements CatalogRepository {
|
||||
.where("status", "=", "running")
|
||||
.returningAll()
|
||||
.execute();
|
||||
return rows.map(serializeSensitiveDataSuggestionRun);
|
||||
return rows.map(serializeSensitivityAnalysisRun);
|
||||
}
|
||||
|
||||
async updateSensitiveDataSuggestionRun(
|
||||
async updateSensitivityAnalysisRun(
|
||||
runId: string,
|
||||
update: SensitiveDataSuggestionRunUpdate,
|
||||
): Promise<SensitiveDataSuggestionRun | undefined> {
|
||||
update: SensitivityAnalysisRunUpdate,
|
||||
): Promise<SensitivityAnalysisRun | undefined> {
|
||||
const values: any = { ...update, updatedAt: sql`now()` };
|
||||
const row = await this.db.updateTable("sensitiveDataSuggestionRuns")
|
||||
.set(values)
|
||||
.where("id", "=", runId)
|
||||
.returningAll()
|
||||
.executeTakeFirst();
|
||||
return row ? serializeSensitiveDataSuggestionRun(row) : undefined;
|
||||
return row ? serializeSensitivityAnalysisRun(row) : undefined;
|
||||
}
|
||||
|
||||
async appendSensitiveDataSuggestionEvent(
|
||||
async appendSensitivityAnalysisEvent(
|
||||
runId: string,
|
||||
level: SensitiveDataSuggestionEvent["level"],
|
||||
level: SensitivityAnalysisEvent["level"],
|
||||
message: string,
|
||||
): Promise<SensitiveDataSuggestionEvent> {
|
||||
): Promise<SensitivityAnalysisEvent> {
|
||||
return await this.db.transaction().execute(async (trx) => {
|
||||
const run = await trx.selectFrom("sensitiveDataSuggestionRuns")
|
||||
.select("id")
|
||||
.where("id", "=", runId)
|
||||
.forUpdate()
|
||||
.executeTakeFirst();
|
||||
if (!run) throw new CatalogConflictError("Sensitive Data Suggestion Run does not exist");
|
||||
if (!run) throw new CatalogConflictError("Sensitivity Analysis Run does not exist");
|
||||
const current = await trx.selectFrom("sensitiveDataSuggestionEvents")
|
||||
.select(sql<number>`coalesce(max(sequence), 0)::int`.as("sequence"))
|
||||
.where("runId", "=", runId)
|
||||
@@ -1032,21 +1039,21 @@ export class KyselyCatalogRepository implements CatalogRepository {
|
||||
level,
|
||||
message,
|
||||
}).returningAll().executeTakeFirstOrThrow();
|
||||
return serializeSensitiveDataSuggestionEvent(row);
|
||||
return serializeSensitivityAnalysisEvent(row);
|
||||
});
|
||||
}
|
||||
|
||||
async listSensitiveDataSuggestionEvents(
|
||||
async listSensitivityAnalysisEvents(
|
||||
runId: string,
|
||||
afterSequence = 0,
|
||||
): Promise<SensitiveDataSuggestionEvent[]> {
|
||||
): Promise<SensitivityAnalysisEvent[]> {
|
||||
const rows = await this.db.selectFrom("sensitiveDataSuggestionEvents")
|
||||
.selectAll()
|
||||
.where("runId", "=", runId)
|
||||
.where("sequence", ">", afterSequence)
|
||||
.orderBy("sequence")
|
||||
.execute();
|
||||
return rows.map(serializeSensitiveDataSuggestionEvent);
|
||||
return rows.map(serializeSensitivityAnalysisEvent);
|
||||
}
|
||||
|
||||
async listRelationships(databaseId: string): Promise<CatalogPhysicalRelationship[]> {
|
||||
@@ -1792,13 +1799,13 @@ export class UnavailableCatalogRepository implements CatalogRepository {
|
||||
async updateDescriptionGenerationRun(): Promise<DescriptionGenerationRun | undefined> { return this.fail(); }
|
||||
async appendDescriptionGenerationEvent(): Promise<DescriptionGenerationEvent> { return this.fail(); }
|
||||
async listDescriptionGenerationEvents(): Promise<DescriptionGenerationEvent[]> { return this.fail(); }
|
||||
async createSensitiveDataSuggestionRun(): Promise<SensitiveDataSuggestionRun> { return this.fail(); }
|
||||
async getSensitiveDataSuggestionRun(): Promise<SensitiveDataSuggestionRun | undefined> { return this.fail(); }
|
||||
async listSensitiveDataSuggestionRuns(): Promise<SensitiveDataSuggestionRun[]> { return this.fail(); }
|
||||
async interruptActiveSensitiveDataSuggestionRuns(): Promise<SensitiveDataSuggestionRun[]> { return this.fail(); }
|
||||
async updateSensitiveDataSuggestionRun(): Promise<SensitiveDataSuggestionRun | undefined> { return this.fail(); }
|
||||
async appendSensitiveDataSuggestionEvent(): Promise<SensitiveDataSuggestionEvent> { return this.fail(); }
|
||||
async listSensitiveDataSuggestionEvents(): Promise<SensitiveDataSuggestionEvent[]> { return this.fail(); }
|
||||
async createSensitivityAnalysisRun(): Promise<SensitivityAnalysisRun> { return this.fail(); }
|
||||
async getSensitivityAnalysisRun(): Promise<SensitivityAnalysisRun | undefined> { return this.fail(); }
|
||||
async listSensitivityAnalysisRuns(): Promise<SensitivityAnalysisRun[]> { return this.fail(); }
|
||||
async interruptActiveSensitivityAnalysisRuns(): Promise<SensitivityAnalysisRun[]> { return this.fail(); }
|
||||
async updateSensitivityAnalysisRun(): Promise<SensitivityAnalysisRun | undefined> { return this.fail(); }
|
||||
async appendSensitivityAnalysisEvent(): Promise<SensitivityAnalysisEvent> { return this.fail(); }
|
||||
async listSensitivityAnalysisEvents(): Promise<SensitivityAnalysisEvent[]> { return this.fail(); }
|
||||
async listRelationships(): Promise<CatalogPhysicalRelationship[]> { return this.fail(); }
|
||||
async listLogicalRelationships(): Promise<CatalogLogicalRelationship[]> { return this.fail(); }
|
||||
async getLogicalRelationshipContext(): Promise<CatalogLogicalRelationshipContext | undefined> { return this.fail(); }
|
||||
|
||||
@@ -1,254 +0,0 @@
|
||||
import { z } from "zod";
|
||||
import type { MetadataGenerationModels } from "./metadata-generation-models.js";
|
||||
import type { ModelCompleter, ModelCompletionMessage, ModelCompletionResult, ModelCompletionUsage } from "./model-completer.js";
|
||||
import type {
|
||||
CatalogColumn,
|
||||
CatalogRepository,
|
||||
CatalogTable,
|
||||
SensitiveDataSuggestionScope,
|
||||
} from "./types.js";
|
||||
|
||||
export type { SensitiveDataSuggestionScope } from "./types.js";
|
||||
|
||||
// The helper accepts at most 64 KiB per message. Keep the same safety margin used by
|
||||
// Description Generation so UTF-8 structural metadata never reaches that hard limit.
|
||||
const MAX_USER_MESSAGE_BYTES = 60 * 1024;
|
||||
// Preserve ThothAI's proven completion granularity: small batches keep generation time and
|
||||
// structured-output accuracy predictable even when the helper byte limit would allow much more.
|
||||
const MAX_COLUMNS_PER_BATCH = 10;
|
||||
const responseSchema = z.object({
|
||||
suggestions: z.array(z.object({
|
||||
columnId: z.uuid(),
|
||||
sensitive: z.boolean(),
|
||||
}).strict()),
|
||||
}).strict();
|
||||
|
||||
interface StructuralColumn {
|
||||
columnId: string;
|
||||
tableId: string;
|
||||
table: string;
|
||||
column: string;
|
||||
dataType: string;
|
||||
nullable: boolean;
|
||||
primaryKey: boolean;
|
||||
foreignKey: boolean;
|
||||
version: number;
|
||||
currentSensitive: boolean;
|
||||
}
|
||||
|
||||
export interface SensitiveDataSuggestion {
|
||||
columnId: string;
|
||||
tableId: string;
|
||||
tableName: string;
|
||||
columnName: string;
|
||||
version: number;
|
||||
currentSensitive: boolean;
|
||||
sensitive: boolean;
|
||||
}
|
||||
|
||||
export class SensitiveDataSuggestionTargetNotFoundError extends Error {
|
||||
constructor(readonly target: "database" | "table" | "column") {
|
||||
super(`${target} not found`);
|
||||
this.name = "SensitiveDataSuggestionTargetNotFoundError";
|
||||
}
|
||||
}
|
||||
|
||||
export class SensitiveDataSuggestionDuplicateTargetIdsError extends Error {
|
||||
constructor() {
|
||||
super("sensitive-data suggestion target IDs must be unique");
|
||||
this.name = "SensitiveDataSuggestionDuplicateTargetIdsError";
|
||||
}
|
||||
}
|
||||
|
||||
export class SensitiveDataSuggestionNoEligibleColumnsError extends Error {
|
||||
constructor(readonly scope: SensitiveDataSuggestionScope) {
|
||||
super("selected scope has no catalog columns");
|
||||
this.name = "SensitiveDataSuggestionNoEligibleColumnsError";
|
||||
}
|
||||
}
|
||||
|
||||
export class SensitiveDataSuggestionPayloadTooLargeError extends Error {
|
||||
constructor() {
|
||||
super("sensitive-data suggestion structural metadata is too large");
|
||||
this.name = "SensitiveDataSuggestionPayloadTooLargeError";
|
||||
}
|
||||
}
|
||||
|
||||
export class SensitiveDataSuggestionInvalidResponseError extends Error {
|
||||
constructor() {
|
||||
super("sensitive-data suggestion response is invalid");
|
||||
this.name = "SensitiveDataSuggestionInvalidResponseError";
|
||||
}
|
||||
}
|
||||
|
||||
function userContent(
|
||||
database: { databaseName: string; schema: string },
|
||||
columns: readonly StructuralColumn[],
|
||||
): string {
|
||||
return JSON.stringify({
|
||||
database: database.databaseName,
|
||||
schema: database.schema,
|
||||
columns: columns.map((column) => ({
|
||||
columnId: column.columnId,
|
||||
table: column.table,
|
||||
column: column.column,
|
||||
dataType: column.dataType,
|
||||
nullable: column.nullable,
|
||||
primaryKey: column.primaryKey,
|
||||
foreignKey: column.foreignKey,
|
||||
})),
|
||||
});
|
||||
}
|
||||
|
||||
function batchesFor(
|
||||
database: { databaseName: string; schema: string },
|
||||
columns: readonly StructuralColumn[],
|
||||
): StructuralColumn[][] {
|
||||
const batches: StructuralColumn[][] = [];
|
||||
let current: StructuralColumn[] = [];
|
||||
for (const column of columns) {
|
||||
if (current.length === MAX_COLUMNS_PER_BATCH) {
|
||||
batches.push(current);
|
||||
current = [];
|
||||
}
|
||||
const candidate = [...current, column];
|
||||
if (Buffer.byteLength(userContent(database, candidate), "utf8") <= MAX_USER_MESSAGE_BYTES) {
|
||||
current = candidate;
|
||||
continue;
|
||||
}
|
||||
if (current.length === 0) throw new SensitiveDataSuggestionPayloadTooLargeError();
|
||||
batches.push(current);
|
||||
current = [column];
|
||||
if (Buffer.byteLength(userContent(database, current), "utf8") > MAX_USER_MESSAGE_BYTES) {
|
||||
throw new SensitiveDataSuggestionPayloadTooLargeError();
|
||||
}
|
||||
}
|
||||
if (current.length > 0) batches.push(current);
|
||||
return batches;
|
||||
}
|
||||
|
||||
function structuralColumn(table: CatalogTable, column: CatalogColumn): StructuralColumn {
|
||||
return {
|
||||
columnId: column.id,
|
||||
tableId: table.id,
|
||||
table: table.name,
|
||||
column: column.name,
|
||||
dataType: column.dataType,
|
||||
nullable: column.isNullable,
|
||||
primaryKey: column.isPrimaryKey,
|
||||
foreignKey: column.isForeignKey,
|
||||
version: column.version,
|
||||
currentSensitive: column.sensitive,
|
||||
};
|
||||
}
|
||||
|
||||
const systemMessage: ModelCompletionMessage = {
|
||||
role: "system",
|
||||
content: [
|
||||
"Classify whether each database column is likely to contain sensitive source values.",
|
||||
"Use only the supplied structural metadata. Return strict JSON with this exact shape:",
|
||||
'{"suggestions":[{"columnId":"uuid","sensitive":true}]}',
|
||||
"Return every supplied column exactly once. Do not add explanations or markdown.",
|
||||
].join("\n"),
|
||||
};
|
||||
|
||||
export class SensitiveDataSuggester {
|
||||
constructor(
|
||||
private readonly repository: CatalogRepository,
|
||||
private readonly models: MetadataGenerationModels,
|
||||
private readonly completer: ModelCompleter,
|
||||
) {}
|
||||
|
||||
private async selectColumns(
|
||||
databaseId: string,
|
||||
scope: SensitiveDataSuggestionScope,
|
||||
targetIds: readonly string[],
|
||||
): Promise<StructuralColumn[]> {
|
||||
if (new Set(targetIds).size !== targetIds.length) {
|
||||
throw new SensitiveDataSuggestionDuplicateTargetIdsError();
|
||||
}
|
||||
const tables = await this.repository.listTables(databaseId);
|
||||
const tableIds = new Set(targetIds);
|
||||
const selectedTables = scope === "selected_tables"
|
||||
? tables.filter((table) => tableIds.has(table.id))
|
||||
: tables;
|
||||
if (scope === "selected_tables" && selectedTables.length !== targetIds.length) {
|
||||
throw new SensitiveDataSuggestionTargetNotFoundError("table");
|
||||
}
|
||||
|
||||
const columns = (await Promise.all(selectedTables.map(async (table) => (
|
||||
(await this.repository.listColumns(databaseId, table.id)).map((column) => (
|
||||
structuralColumn(table, column)
|
||||
))
|
||||
)))).flat();
|
||||
const columnIds = new Set(targetIds);
|
||||
const selectedColumns = scope === "selected_columns"
|
||||
? columns.filter((column) => columnIds.has(column.columnId))
|
||||
: columns;
|
||||
if (scope === "selected_columns" && selectedColumns.length !== targetIds.length) {
|
||||
throw new SensitiveDataSuggestionTargetNotFoundError("column");
|
||||
}
|
||||
if (selectedColumns.length === 0) {
|
||||
throw new SensitiveDataSuggestionNoEligibleColumnsError(scope);
|
||||
}
|
||||
return selectedColumns;
|
||||
}
|
||||
|
||||
async suggest(
|
||||
databaseId: string,
|
||||
modelId: string,
|
||||
scope: SensitiveDataSuggestionScope,
|
||||
targetIds: readonly string[],
|
||||
signal: AbortSignal,
|
||||
onPrepared?: (total: number) => void | Promise<void>,
|
||||
onProgress?: (processed: number, suggestions: readonly SensitiveDataSuggestion[]) => void | Promise<void>,
|
||||
onUsage?: (usage: ModelCompletionUsage) => void | Promise<void>,
|
||||
): Promise<readonly SensitiveDataSuggestion[]> {
|
||||
const database = await this.repository.get(databaseId);
|
||||
if (!database) throw new SensitiveDataSuggestionTargetNotFoundError("database");
|
||||
const columns = await this.selectColumns(databaseId, scope, targetIds);
|
||||
await onPrepared?.(columns.length);
|
||||
const model = this.models.resolve(modelId);
|
||||
const suggestions: SensitiveDataSuggestion[] = [];
|
||||
|
||||
for (const batch of batchesFor(database, columns)) {
|
||||
let received: Map<string, { columnId: string; sensitive: boolean }> | undefined;
|
||||
for (let attempt = 0; attempt < 2 && !received; attempt += 1) {
|
||||
const completion = await this.completer.complete({
|
||||
model,
|
||||
signal,
|
||||
messages: [systemMessage, { role: "user", content: userContent(database, batch) }],
|
||||
});
|
||||
const result: ModelCompletionResult = typeof completion === "string"
|
||||
? { content: completion, usage: { input: 0, cacheRead: 0, output: 0 } }
|
||||
: completion;
|
||||
await onUsage?.(result.usage);
|
||||
const content = result.content;
|
||||
try {
|
||||
const parsed = responseSchema.parse(JSON.parse(content));
|
||||
const expected = new Set(batch.map((column) => column.columnId));
|
||||
const candidate = new Map(parsed.suggestions.map((suggestion) => [suggestion.columnId, suggestion]));
|
||||
if (candidate.size !== parsed.suggestions.length
|
||||
|| candidate.size !== expected.size
|
||||
|| [...candidate.keys()].some((columnId) => !expected.has(columnId))) {
|
||||
throw new SensitiveDataSuggestionInvalidResponseError();
|
||||
}
|
||||
received = candidate;
|
||||
} catch {
|
||||
if (attempt === 1) throw new SensitiveDataSuggestionInvalidResponseError();
|
||||
}
|
||||
}
|
||||
suggestions.push(...batch.map((column) => ({
|
||||
columnId: column.columnId,
|
||||
tableId: column.tableId,
|
||||
tableName: column.table,
|
||||
columnName: column.column,
|
||||
version: column.version,
|
||||
currentSensitive: column.currentSensitive,
|
||||
sensitive: received!.get(column.columnId)!.sensitive,
|
||||
})));
|
||||
await onProgress?.(suggestions.length, suggestions.slice(-batch.length));
|
||||
}
|
||||
return suggestions;
|
||||
}
|
||||
}
|
||||
@@ -1,136 +0,0 @@
|
||||
import type {
|
||||
SensitiveDataSuggestion,
|
||||
} from "./sensitive-data-suggester.js";
|
||||
import {
|
||||
SensitiveDataSuggester,
|
||||
SensitiveDataSuggestionTargetNotFoundError,
|
||||
} from "./sensitive-data-suggester.js";
|
||||
import type {
|
||||
CatalogRepository,
|
||||
SensitiveDataSuggestionRun,
|
||||
SensitiveDataSuggestionScope,
|
||||
} from "./types.js";
|
||||
import type { ModelCompletionUsage } from "./model-completer.js";
|
||||
|
||||
const interruptedMessage = "Sensitive-field suggestion generation was interrupted by backend restart.";
|
||||
const failedMessage = "Sensitive-field suggestion generation failed.";
|
||||
|
||||
export interface SensitiveDataSuggestionRunResult {
|
||||
suggestions: readonly SensitiveDataSuggestion[];
|
||||
run: SensitiveDataSuggestionRun;
|
||||
}
|
||||
|
||||
export class SensitiveDataSuggestionRunner {
|
||||
constructor(
|
||||
private readonly repository: CatalogRepository,
|
||||
private readonly suggester: SensitiveDataSuggester,
|
||||
) {}
|
||||
|
||||
async initialize(): Promise<void> {
|
||||
if (!(await this.repository.available())) return;
|
||||
const interrupted = await this.repository.interruptActiveSensitiveDataSuggestionRuns(
|
||||
interruptedMessage,
|
||||
);
|
||||
for (const run of interrupted) {
|
||||
await this.repository.appendSensitiveDataSuggestionEvent(
|
||||
run.id,
|
||||
"warning",
|
||||
interruptedMessage,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
async run(
|
||||
databaseId: string,
|
||||
modelId: string,
|
||||
scope: SensitiveDataSuggestionScope,
|
||||
targetIds: readonly string[],
|
||||
signal: AbortSignal,
|
||||
): Promise<SensitiveDataSuggestionRunResult> {
|
||||
if (!(await this.repository.get(databaseId))) {
|
||||
throw new SensitiveDataSuggestionTargetNotFoundError("database");
|
||||
}
|
||||
const started = await this.repository.createSensitiveDataSuggestionRun(
|
||||
databaseId,
|
||||
scope,
|
||||
modelId,
|
||||
);
|
||||
|
||||
try {
|
||||
await this.repository.appendSensitiveDataSuggestionEvent(
|
||||
started.id,
|
||||
"info",
|
||||
"Sensitive-field suggestion generation started.",
|
||||
);
|
||||
const suggestions = await this.suggester.suggest(
|
||||
databaseId,
|
||||
modelId,
|
||||
scope,
|
||||
targetIds,
|
||||
signal,
|
||||
async (total) => {
|
||||
const prepared = await this.repository.updateSensitiveDataSuggestionRun(started.id, {
|
||||
total,
|
||||
});
|
||||
if (!prepared) throw new Error("Sensitive Data Suggestion Run disappeared");
|
||||
},
|
||||
async (processed, batch) => {
|
||||
const suggestedSensitive = batch.filter((suggestion) => suggestion.sensitive).length;
|
||||
const suggestedNonSensitive = batch.length - suggestedSensitive;
|
||||
const current = await this.repository.getSensitiveDataSuggestionRun(started.id);
|
||||
if (!current) throw new Error("Sensitive Data Suggestion Run disappeared");
|
||||
const progress = await this.repository.updateSensitiveDataSuggestionRun(started.id, {
|
||||
suggestedSensitive: current.suggestedSensitive + suggestedSensitive,
|
||||
suggestedNonSensitive: current.suggestedNonSensitive + suggestedNonSensitive,
|
||||
});
|
||||
if (!progress) throw new Error("Sensitive Data Suggestion Run disappeared");
|
||||
await this.repository.appendSensitiveDataSuggestionEvent(
|
||||
started.id,
|
||||
"info",
|
||||
`Classified ${processed} of ${progress.total} columns.`,
|
||||
);
|
||||
},
|
||||
async (usage: ModelCompletionUsage) => {
|
||||
const current = await this.repository.getSensitiveDataSuggestionRun(started.id);
|
||||
if (!current) throw new Error("Sensitive Data Suggestion Run disappeared");
|
||||
await this.repository.updateSensitiveDataSuggestionRun(started.id, {
|
||||
inputTokens: current.inputTokens + usage.input,
|
||||
cacheReadTokens: current.cacheReadTokens + usage.cacheRead,
|
||||
outputTokens: current.outputTokens + usage.output,
|
||||
});
|
||||
},
|
||||
);
|
||||
const suggestedSensitive = suggestions.filter((suggestion) => suggestion.sensitive).length;
|
||||
const suggestedNonSensitive = suggestions.length - suggestedSensitive;
|
||||
await this.repository.appendSensitiveDataSuggestionEvent(
|
||||
started.id,
|
||||
"info",
|
||||
`Sensitive-field suggestion generation completed for ${suggestions.length} column${
|
||||
suggestions.length === 1 ? "" : "s"
|
||||
}.`,
|
||||
);
|
||||
const completed = await this.repository.updateSensitiveDataSuggestionRun(started.id, {
|
||||
status: "completed",
|
||||
total: suggestions.length,
|
||||
suggestedSensitive,
|
||||
suggestedNonSensitive,
|
||||
finishedAt: new Date().toISOString(),
|
||||
errorSummary: null,
|
||||
});
|
||||
if (!completed) throw new Error("Sensitive Data Suggestion Run disappeared");
|
||||
return { suggestions, run: completed };
|
||||
} catch (error) {
|
||||
await this.repository.updateSensitiveDataSuggestionRun(started.id, {
|
||||
status: "failed",
|
||||
finishedAt: new Date().toISOString(),
|
||||
errorSummary: failedMessage,
|
||||
}).catch(() => undefined);
|
||||
await this.repository.appendSensitiveDataSuggestionEvent(
|
||||
started.id,
|
||||
"error",
|
||||
failedMessage,
|
||||
).catch(() => undefined);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,172 @@
|
||||
import type {
|
||||
SensitivityReviewItem,
|
||||
} from "./sensitivity-analysis-service.js";
|
||||
import {
|
||||
SENSITIVITY_POLICY_VERSION,
|
||||
SensitivityAnalysisInterruptedError,
|
||||
SensitivityAnalysisService,
|
||||
SensitivityAnalysisTargetNotFoundError,
|
||||
} from "./sensitivity-analysis-service.js";
|
||||
import type {
|
||||
CatalogRepository,
|
||||
SensitivityAnalysisRun,
|
||||
SensitivityAnalysisScope,
|
||||
} from "./types.js";
|
||||
|
||||
const interruptedMessage = "Local sensitivity analysis was interrupted by backend restart.";
|
||||
const deadlineMessage = "Local sensitivity analysis reached its time limit.";
|
||||
const failedMessage = "Local sensitivity analysis failed.";
|
||||
|
||||
function ensureActive(signal: AbortSignal): void {
|
||||
if (signal.aborted) throw new SensitivityAnalysisInterruptedError();
|
||||
}
|
||||
|
||||
export interface SensitivityAnalysisRunResult {
|
||||
suggestions: readonly SensitivityReviewItem[];
|
||||
run: SensitivityAnalysisRun;
|
||||
}
|
||||
|
||||
export class SensitivityAnalysisRunner {
|
||||
constructor(
|
||||
private readonly repository: CatalogRepository,
|
||||
private readonly analysis: SensitivityAnalysisService,
|
||||
) {}
|
||||
|
||||
async initialize(): Promise<void> {
|
||||
if (!(await this.repository.available())) return;
|
||||
const interrupted = await this.repository.interruptActiveSensitivityAnalysisRuns(
|
||||
interruptedMessage,
|
||||
);
|
||||
for (const run of interrupted) {
|
||||
await this.repository.appendSensitivityAnalysisEvent(
|
||||
run.id,
|
||||
"warning",
|
||||
interruptedMessage,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
async run(
|
||||
databaseId: string,
|
||||
scope: SensitivityAnalysisScope,
|
||||
targetIds: readonly string[],
|
||||
signal: AbortSignal,
|
||||
): Promise<SensitivityAnalysisRunResult> {
|
||||
ensureActive(signal);
|
||||
const database = await this.repository.get(databaseId);
|
||||
ensureActive(signal);
|
||||
if (!database) {
|
||||
throw new SensitivityAnalysisTargetNotFoundError("database");
|
||||
}
|
||||
ensureActive(signal);
|
||||
const started = await this.repository.createSensitivityAnalysisRun(
|
||||
databaseId,
|
||||
scope,
|
||||
{ engine: "local", policyVersion: SENSITIVITY_POLICY_VERSION },
|
||||
);
|
||||
let preparedTotal = 0;
|
||||
let processedSensitive = 0;
|
||||
let processedNonSensitive = 0;
|
||||
|
||||
try {
|
||||
ensureActive(signal);
|
||||
await this.repository.appendSensitivityAnalysisEvent(
|
||||
started.id,
|
||||
"info",
|
||||
"Local sensitivity analysis started.",
|
||||
);
|
||||
ensureActive(signal);
|
||||
const suggestions = await this.analysis.analyze(
|
||||
databaseId,
|
||||
scope,
|
||||
targetIds,
|
||||
signal,
|
||||
async (total) => {
|
||||
ensureActive(signal);
|
||||
preparedTotal = total;
|
||||
const prepared = await this.repository.updateSensitivityAnalysisRun(started.id, {
|
||||
total,
|
||||
});
|
||||
ensureActive(signal);
|
||||
if (!prepared) throw new Error("Sensitivity Analysis Run disappeared");
|
||||
},
|
||||
async (processed, batch) => {
|
||||
ensureActive(signal);
|
||||
const suggestedSensitive = batch.filter(
|
||||
(suggestion) => suggestion.assessment === "sensitive",
|
||||
).length;
|
||||
const suggestedNonSensitive = batch.filter(
|
||||
(suggestion) => suggestion.assessment === "non_sensitive",
|
||||
).length;
|
||||
const unknown = batch.filter((suggestion) => suggestion.assessment === "unknown").length;
|
||||
const current = await this.repository.getSensitivityAnalysisRun(started.id);
|
||||
ensureActive(signal);
|
||||
if (!current) throw new Error("Sensitivity Analysis Run disappeared");
|
||||
const progress = await this.repository.updateSensitivityAnalysisRun(started.id, {
|
||||
suggestedSensitive: current.suggestedSensitive + suggestedSensitive,
|
||||
suggestedNonSensitive: current.suggestedNonSensitive + suggestedNonSensitive,
|
||||
unknown: current.unknown + unknown,
|
||||
});
|
||||
if (!progress) throw new Error("Sensitivity Analysis Run disappeared");
|
||||
processedSensitive += suggestedSensitive;
|
||||
processedNonSensitive += suggestedNonSensitive;
|
||||
ensureActive(signal);
|
||||
await this.repository.appendSensitivityAnalysisEvent(
|
||||
started.id,
|
||||
"info",
|
||||
`Assessed ${processed} of ${progress.total} columns locally.`,
|
||||
);
|
||||
ensureActive(signal);
|
||||
},
|
||||
);
|
||||
ensureActive(signal);
|
||||
const suggestedSensitive = suggestions.filter(
|
||||
(suggestion) => suggestion.assessment === "sensitive",
|
||||
).length;
|
||||
const suggestedNonSensitive = suggestions.filter(
|
||||
(suggestion) => suggestion.assessment === "non_sensitive",
|
||||
).length;
|
||||
const unknown = suggestions.filter((suggestion) => suggestion.assessment === "unknown").length;
|
||||
await this.repository.appendSensitivityAnalysisEvent(
|
||||
started.id,
|
||||
"info",
|
||||
`Local sensitivity analysis completed for ${suggestions.length} column${
|
||||
suggestions.length === 1 ? "" : "s"
|
||||
}.`,
|
||||
);
|
||||
ensureActive(signal);
|
||||
const completed = await this.repository.updateSensitivityAnalysisRun(started.id, {
|
||||
status: "completed",
|
||||
total: suggestions.length,
|
||||
suggestedSensitive,
|
||||
suggestedNonSensitive,
|
||||
unknown,
|
||||
finishedAt: new Date().toISOString(),
|
||||
errorSummary: null,
|
||||
});
|
||||
ensureActive(signal);
|
||||
if (!completed) throw new Error("Sensitivity Analysis Run disappeared");
|
||||
return { suggestions, run: completed };
|
||||
} catch (error) {
|
||||
const interrupted = signal.aborted || error instanceof SensitivityAnalysisInterruptedError;
|
||||
const message = interrupted ? deadlineMessage : failedMessage;
|
||||
await this.repository.updateSensitivityAnalysisRun(started.id, {
|
||||
status: interrupted ? "interrupted" : "failed",
|
||||
...(interrupted ? {
|
||||
total: preparedTotal,
|
||||
suggestedSensitive: processedSensitive,
|
||||
suggestedNonSensitive: processedNonSensitive,
|
||||
unknown: Math.max(0, preparedTotal - processedSensitive - processedNonSensitive),
|
||||
} : {}),
|
||||
finishedAt: new Date().toISOString(),
|
||||
errorSummary: message,
|
||||
}).catch(() => undefined);
|
||||
await this.repository.appendSensitivityAnalysisEvent(
|
||||
started.id,
|
||||
interrupted ? "warning" : "error",
|
||||
message,
|
||||
).catch(() => undefined);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,175 @@
|
||||
import type {
|
||||
SensitivityClassifier,
|
||||
SensitivityColumnAssessment,
|
||||
SensitivityEvidence,
|
||||
SensitivityNerBudget,
|
||||
} from "./sensitivity-classifier.js";
|
||||
import type {
|
||||
CatalogColumn,
|
||||
CatalogRepository,
|
||||
CatalogTable,
|
||||
SensitivityAnalysisScope,
|
||||
} from "./types.js";
|
||||
|
||||
export type { SensitivityAnalysisScope } from "./types.js";
|
||||
export const SENSITIVITY_POLICY_VERSION = "sensitivity-v1";
|
||||
|
||||
interface SelectedColumn {
|
||||
table: CatalogTable;
|
||||
column: CatalogColumn;
|
||||
}
|
||||
|
||||
export interface SensitivityReviewItem {
|
||||
columnId: string;
|
||||
tableId: string;
|
||||
tableName: string;
|
||||
columnName: string;
|
||||
version: number;
|
||||
currentSensitive: boolean;
|
||||
sensitive: boolean;
|
||||
assessment: SensitivityColumnAssessment["assessment"];
|
||||
evidence: readonly SensitivityEvidence[];
|
||||
observedValues: number;
|
||||
}
|
||||
|
||||
export class SensitivityAnalysisTargetNotFoundError extends Error {
|
||||
constructor(readonly target: "database" | "table" | "column") {
|
||||
super(`${target} not found`);
|
||||
this.name = "SensitivityAnalysisTargetNotFoundError";
|
||||
}
|
||||
}
|
||||
|
||||
export class SensitivityAnalysisDuplicateTargetIdsError extends Error {
|
||||
constructor() {
|
||||
super("sensitivity analysis target IDs must be unique");
|
||||
this.name = "SensitivityAnalysisDuplicateTargetIdsError";
|
||||
}
|
||||
}
|
||||
|
||||
export class SensitivityAnalysisInterruptedError extends Error {
|
||||
constructor() {
|
||||
super("sensitivity analysis deadline exceeded");
|
||||
this.name = "SensitivityAnalysisInterruptedError";
|
||||
}
|
||||
}
|
||||
|
||||
function ensureActive(signal: AbortSignal): void {
|
||||
if (signal.aborted) throw new SensitivityAnalysisInterruptedError();
|
||||
}
|
||||
|
||||
export class SensitivityAnalysisNoEligibleColumnsError extends Error {
|
||||
constructor(readonly scope: SensitivityAnalysisScope) {
|
||||
super("selected scope has no catalog columns");
|
||||
this.name = "SensitivityAnalysisNoEligibleColumnsError";
|
||||
}
|
||||
}
|
||||
|
||||
/** Selection and table orchestration around the single SensitivityClassifier decision module. */
|
||||
export class SensitivityAnalysisService {
|
||||
constructor(
|
||||
private readonly repository: CatalogRepository,
|
||||
private readonly classifier: SensitivityClassifier,
|
||||
private readonly options: { runBudgetMs?: number; nerBudgetMs?: number; now?: () => number } = {},
|
||||
) {}
|
||||
|
||||
private async selectColumns(
|
||||
databaseId: string,
|
||||
scope: SensitivityAnalysisScope,
|
||||
targetIds: readonly string[],
|
||||
signal: AbortSignal,
|
||||
): Promise<readonly SelectedColumn[]> {
|
||||
ensureActive(signal);
|
||||
if (new Set(targetIds).size !== targetIds.length) {
|
||||
throw new SensitivityAnalysisDuplicateTargetIdsError();
|
||||
}
|
||||
const tables = await this.repository.listTables(databaseId);
|
||||
ensureActive(signal);
|
||||
const tableIds = new Set(targetIds);
|
||||
const selectedTables = scope === "selected_tables"
|
||||
? tables.filter((table) => tableIds.has(table.id))
|
||||
: tables;
|
||||
if (scope === "selected_tables" && selectedTables.length !== targetIds.length) {
|
||||
throw new SensitivityAnalysisTargetNotFoundError("table");
|
||||
}
|
||||
const columns = (await Promise.all(selectedTables.map(async (table) => (
|
||||
(await this.repository.listColumns(databaseId, table.id)).map((column) => ({ table, column }))
|
||||
)))).flat();
|
||||
ensureActive(signal);
|
||||
const columnIds = new Set(targetIds);
|
||||
const selectedColumns = scope === "selected_columns"
|
||||
? columns.filter(({ column }) => columnIds.has(column.id))
|
||||
: columns;
|
||||
if (scope === "selected_columns" && selectedColumns.length !== targetIds.length) {
|
||||
throw new SensitivityAnalysisTargetNotFoundError("column");
|
||||
}
|
||||
if (selectedColumns.length === 0) {
|
||||
throw new SensitivityAnalysisNoEligibleColumnsError(scope);
|
||||
}
|
||||
return selectedColumns;
|
||||
}
|
||||
|
||||
async analyze(
|
||||
databaseId: string,
|
||||
scope: SensitivityAnalysisScope,
|
||||
targetIds: readonly string[],
|
||||
signal: AbortSignal,
|
||||
onPrepared?: (total: number) => void | Promise<void>,
|
||||
onProgress?: (processed: number, suggestions: readonly SensitivityReviewItem[]) => void | Promise<void>,
|
||||
): Promise<readonly SensitivityReviewItem[]> {
|
||||
const now = this.options.now ?? Date.now;
|
||||
const deadline = now() + (this.options.runBudgetMs ?? 60_000);
|
||||
const configuredNerBudget = this.options.nerBudgetMs ?? 10_000;
|
||||
const nerBudget: SensitivityNerBudget = {
|
||||
remainingMs: Number.isFinite(configuredNerBudget) && configuredNerBudget >= 0
|
||||
? configuredNerBudget
|
||||
: 10_000,
|
||||
};
|
||||
ensureActive(signal);
|
||||
const database = await this.repository.get(databaseId);
|
||||
ensureActive(signal);
|
||||
if (!database) throw new SensitivityAnalysisTargetNotFoundError("database");
|
||||
const selected = await this.selectColumns(databaseId, scope, targetIds, signal);
|
||||
await onPrepared?.(selected.length);
|
||||
ensureActive(signal);
|
||||
const byTable = new Map<string, SelectedColumn[]>();
|
||||
for (const item of selected) {
|
||||
const items = byTable.get(item.table.id) ?? [];
|
||||
items.push(item);
|
||||
byTable.set(item.table.id, items);
|
||||
}
|
||||
const suggestions: SensitivityReviewItem[] = [];
|
||||
for (const items of byTable.values()) {
|
||||
ensureActive(signal);
|
||||
const first = items[0]!;
|
||||
const assessments = await this.classifier.assessTable({
|
||||
database,
|
||||
table: first.table,
|
||||
columns: items.map(({ column }) => column),
|
||||
}, signal, deadline, nerBudget);
|
||||
ensureActive(signal);
|
||||
const assessmentById = new Map(assessments.map((assessment) => [
|
||||
assessment.columnId,
|
||||
assessment,
|
||||
]));
|
||||
const batch = items.map(({ table, column }) => {
|
||||
const assessment = assessmentById.get(column.id)!;
|
||||
return {
|
||||
columnId: column.id,
|
||||
tableId: table.id,
|
||||
tableName: table.name,
|
||||
columnName: column.name,
|
||||
version: column.version,
|
||||
currentSensitive: column.sensitive,
|
||||
sensitive: assessment.proposedSensitive,
|
||||
assessment: assessment.assessment,
|
||||
evidence: assessment.evidence,
|
||||
observedValues: assessment.observedValues,
|
||||
};
|
||||
});
|
||||
suggestions.push(...batch);
|
||||
await onProgress?.(suggestions.length, batch);
|
||||
ensureActive(signal);
|
||||
}
|
||||
return suggestions;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,441 @@
|
||||
import { CatalogConnectorError, type CatalogColumn, type CatalogTable, type WorkspaceDatabase } from "./types.js";
|
||||
import { findPhoneNumbersInText } from "libphonenumber-js/max";
|
||||
import validator from "validator";
|
||||
|
||||
export type SensitivityAssessment = "sensitive" | "non_sensitive" | "unknown";
|
||||
|
||||
export interface SensitivityEvidence {
|
||||
kind: "metadata" | "content" | "length" | "ner" | "coverage";
|
||||
ruleId: string;
|
||||
label?: string;
|
||||
confidence?: number;
|
||||
}
|
||||
|
||||
export interface SensitivityValueObservation {
|
||||
columnId: string;
|
||||
value: string | null;
|
||||
characterLength: number | null;
|
||||
}
|
||||
|
||||
export interface SensitivityScanCoverage {
|
||||
kind: "complete" | "sampled" | "unavailable";
|
||||
observedRows: number;
|
||||
}
|
||||
|
||||
export interface SensitivityTableScan {
|
||||
batches: readonly (readonly SensitivityValueObservation[])[];
|
||||
coverage: SensitivityScanCoverage;
|
||||
}
|
||||
|
||||
export interface SensitivityScanRequest {
|
||||
database: WorkspaceDatabase;
|
||||
table: CatalogTable;
|
||||
columns: readonly CatalogColumn[];
|
||||
fullScanBudgetMs: number;
|
||||
deadline: number;
|
||||
}
|
||||
|
||||
export interface SensitivityValueSource {
|
||||
scanTable(
|
||||
request: SensitivityScanRequest,
|
||||
consume: (batch: readonly SensitivityValueObservation[]) => void | Promise<void>,
|
||||
signal: AbortSignal,
|
||||
): Promise<SensitivityScanCoverage>;
|
||||
}
|
||||
|
||||
export interface LocalNerCandidate {
|
||||
columnId: string;
|
||||
text: string;
|
||||
}
|
||||
|
||||
export interface LocalNerEvidence {
|
||||
columnId: string;
|
||||
label: string;
|
||||
confidence: number;
|
||||
}
|
||||
|
||||
export interface SensitivityNerBudget {
|
||||
remainingMs: number;
|
||||
}
|
||||
|
||||
/** Optional local detector. It returns evidence only; it never decides a column assessment. */
|
||||
export interface LocalNerDetector {
|
||||
warmup?(): Promise<void>;
|
||||
isReady?(): boolean;
|
||||
detect(
|
||||
candidates: readonly LocalNerCandidate[],
|
||||
signal: AbortSignal,
|
||||
deadline: number,
|
||||
): Promise<readonly LocalNerEvidence[]>;
|
||||
close?(): Promise<void>;
|
||||
}
|
||||
|
||||
export interface SensitivityColumnAssessment {
|
||||
columnId: string;
|
||||
assessment: SensitivityAssessment;
|
||||
proposedSensitive: boolean;
|
||||
evidence: readonly SensitivityEvidence[];
|
||||
observedValues: number;
|
||||
}
|
||||
|
||||
export interface SensitivityTableTarget {
|
||||
database: WorkspaceDatabase;
|
||||
table: CatalogTable;
|
||||
columns: readonly CatalogColumn[];
|
||||
}
|
||||
|
||||
const EMAIL = /(?<![\p{L}\p{N}._%+-])[\p{L}\p{N}._%+-]+@[\p{L}\p{N}.-]+\.[\p{L}]{2,63}(?![\p{L}\p{N}._%+-])/giu;
|
||||
const DIRECT_IDENTIFIER_NAMES = new Set([
|
||||
"address", "birth_date", "codice_fiscale", "date_of_birth", "dob", "email", "e_mail",
|
||||
"bic", "first_name", "fiscal_code", "full_name", "iban", "indirizzo", "last_name", "mobile",
|
||||
"nome", "passport", "phone", "surname", "swift", "swift_code", "tax_id", "telefono",
|
||||
]);
|
||||
const CREDENTIAL_NAME = /(?:^|_)(?:api_key|credential|password|passwd|private_key|pwd|secret|token)(?:_|$)/u;
|
||||
const HEALTH_NAME = /(?:^|_)(?:anamnesi|clinical|diagnos(?:i|is)|health|medical|patient|patologia|therapy|terapia)(?:_|$)/u;
|
||||
const CLINICAL_TERM = /(?:^|[^\p{L}])(?:allergi[ae]|anamnesi|carcinoma|chemioterapia|diabete|diagnos[ei]|epatite|farmac[io]|gravidanza|hiv|metastasi|neoplasia|patologia|radioterapia|referto|terapia|tumore)(?:$|[^\p{L}])/iu;
|
||||
const UNSUPPORTED_BINARY_TYPE = /(?:^|\s)(?:binary|blob|bytea|image|varbinary)(?:\s|$|\()/iu;
|
||||
const MAX_NER_CANDIDATES_PER_REQUEST = 128;
|
||||
|
||||
function normalizedName(value: string): string {
|
||||
return value.normalize("NFKD")
|
||||
.replace(/[\u0300-\u036f]/g, "")
|
||||
.replace(/([a-z0-9])([A-Z])/g, "$1_$2")
|
||||
.toLocaleLowerCase("en-US")
|
||||
.replace(/[^a-z0-9]+/g, "_")
|
||||
.replace(/^_+|_+$/g, "");
|
||||
}
|
||||
|
||||
function boundedCount(value: number | undefined, fallback: number, maximum: number): number {
|
||||
return value === undefined || !Number.isSafeInteger(value)
|
||||
? fallback
|
||||
: Math.max(1, Math.min(value, maximum));
|
||||
}
|
||||
|
||||
function metadataEvidence(column: CatalogColumn): SensitivityEvidence | undefined {
|
||||
const ruleId = sensitiveNameRule(column.name);
|
||||
return ruleId ? { kind: "metadata", ruleId } : undefined;
|
||||
}
|
||||
|
||||
function sensitiveNameRule(value: string): string | undefined {
|
||||
const name = normalizedName(value);
|
||||
if (DIRECT_IDENTIFIER_NAMES.has(name)) {
|
||||
return "metadata.direct_identifier";
|
||||
}
|
||||
if (CREDENTIAL_NAME.test(name)) {
|
||||
return "metadata.credential";
|
||||
}
|
||||
if (HEALTH_NAME.test(name)) {
|
||||
return "metadata.health";
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
const ITALIAN_FISCAL_CODE = /(?<![A-Z0-9])[A-Z]{6}[0-9LMNPQRSTUV]{2}[ABCDEHLMPRST][0-9LMNPQRSTUV]{2}[A-Z][0-9LMNPQRSTUV]{3}[A-Z](?![A-Z0-9])/giu;
|
||||
const FISCAL_ODD: Record<string, number> = {
|
||||
"0": 1, "1": 0, "2": 5, "3": 7, "4": 9, "5": 13, "6": 15, "7": 17, "8": 19, "9": 21,
|
||||
A: 1, B: 0, C: 5, D: 7, E: 9, F: 13, G: 15, H: 17, I: 19, J: 21,
|
||||
K: 2, L: 4, M: 18, N: 20, O: 11, P: 3, Q: 6, R: 8, S: 12, T: 14,
|
||||
U: 16, V: 10, W: 22, X: 25, Y: 24, Z: 23,
|
||||
};
|
||||
|
||||
function validItalianFiscalCode(candidate: string): boolean {
|
||||
const value = candidate.toUpperCase();
|
||||
if (value.length !== 16) return false;
|
||||
let sum = 0;
|
||||
for (let index = 0; index < 15; index += 1) {
|
||||
const character = value[index]!;
|
||||
if (index % 2 === 0) sum += FISCAL_ODD[character] ?? -1000;
|
||||
else sum += /\d/u.test(character) ? Number(character) : character.charCodeAt(0) - 65;
|
||||
}
|
||||
return String.fromCharCode(65 + (sum % 26)) === value[15];
|
||||
}
|
||||
|
||||
function validIban(candidate: string): boolean {
|
||||
const value = candidate.replace(/\s+/gu, "").toUpperCase();
|
||||
if (!/^[A-Z]{2}\d{2}[A-Z0-9]{11,30}$/u.test(value)) return false;
|
||||
const rearranged = value.slice(4) + value.slice(0, 4);
|
||||
let remainder = 0;
|
||||
for (const character of rearranged) {
|
||||
const digits = /\d/u.test(character) ? character : String(character.charCodeAt(0) - 55);
|
||||
for (const digit of digits) remainder = (remainder * 10 + Number(digit)) % 97;
|
||||
}
|
||||
return remainder === 1;
|
||||
}
|
||||
|
||||
function validPaymentCard(candidate: string): boolean {
|
||||
const digits = candidate.replace(/[ -]/gu, "");
|
||||
if (!/^\d{13,19}$/u.test(digits) || /^(\d)\1+$/u.test(digits)) return false;
|
||||
let sum = 0;
|
||||
let double = false;
|
||||
for (let index = digits.length - 1; index >= 0; index -= 1) {
|
||||
let digit = Number(digits[index]);
|
||||
if (double) {
|
||||
digit *= 2;
|
||||
if (digit > 9) digit -= 9;
|
||||
}
|
||||
sum += digit;
|
||||
double = !double;
|
||||
}
|
||||
return sum % 10 === 0;
|
||||
}
|
||||
|
||||
function jsonHasSensitiveKey(value: string): boolean {
|
||||
const trimmed = value.trim();
|
||||
if (!(trimmed.startsWith("{") || trimmed.startsWith("["))) return false;
|
||||
try {
|
||||
const pending: Array<{ value: unknown; depth: number }> = [{ value: JSON.parse(trimmed), depth: 0 }];
|
||||
let visited = 0;
|
||||
while (pending.length > 0 && visited < 1_000) {
|
||||
const item = pending.pop()!;
|
||||
visited += 1;
|
||||
if (item.depth > 8 || item.value === null || typeof item.value !== "object") continue;
|
||||
if (Array.isArray(item.value)) {
|
||||
for (const child of item.value) pending.push({ value: child, depth: item.depth + 1 });
|
||||
continue;
|
||||
}
|
||||
for (const [key, child] of Object.entries(item.value)) {
|
||||
if (sensitiveNameRule(key)) return true;
|
||||
pending.push({ value: child, depth: item.depth + 1 });
|
||||
}
|
||||
}
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
function contentEvidence(value: string): SensitivityEvidence | undefined {
|
||||
if (/-----BEGIN (?:[A-Z0-9]+ )?PRIVATE KEY-----/u.test(value)) {
|
||||
return { kind: "content", ruleId: "credential.private_key" };
|
||||
}
|
||||
if (/(?:^|[^A-Z0-9])AKIA[A-Z0-9]{16}(?![A-Z0-9])/u.test(value)
|
||||
|| /(?:^|[^A-Za-z0-9_])gh[pousr]_[A-Za-z0-9_]{30,}(?![A-Za-z0-9_])/u.test(value)
|
||||
|| /(?:^|[^A-Za-z0-9_-])eyJ[A-Za-z0-9_-]{5,}\.[A-Za-z0-9_-]{5,}\.[A-Za-z0-9_-]{5,}(?![A-Za-z0-9_-])/u.test(value)) {
|
||||
return { kind: "content", ruleId: "credential.access_key" };
|
||||
}
|
||||
if (/(?:^|[^\p{L}\p{N}_])(?:api[_ -]?key|access[_ -]?token|password|passwd|pwd|secret)\s*[:=]\s*[^\s,;]{4,}/iu.test(value)) {
|
||||
return { kind: "content", ruleId: "credential.key_value" };
|
||||
}
|
||||
if (CLINICAL_TERM.test(value)) return { kind: "content", ruleId: "health.clinical_term" };
|
||||
for (const match of value.matchAll(EMAIL)) {
|
||||
if (validator.isEmail(match[0])) return { kind: "content", ruleId: "pii.email" };
|
||||
}
|
||||
for (const match of value.matchAll(ITALIAN_FISCAL_CODE)) {
|
||||
if (validItalianFiscalCode(match[0])) {
|
||||
return { kind: "content", ruleId: "pii.italian_fiscal_code" };
|
||||
}
|
||||
}
|
||||
for (const match of value.matchAll(/\b(?:passaporto|passport)(?:\s+(?:numero|number|n\.?))?\s*[:#-]?\s*([A-Z0-9]{9})\b/giu)) {
|
||||
if (validator.isPassportNumber(match[1]!, "IT")) {
|
||||
return { kind: "content", ruleId: "pii.passport_number" };
|
||||
}
|
||||
}
|
||||
for (const match of value.matchAll(/\bC[A-Z]\d{5}[A-Z]{2}\b/giu)) {
|
||||
if (validator.isIdentityCard(match[0], "IT")) {
|
||||
return { kind: "content", ruleId: "pii.identity_card" };
|
||||
}
|
||||
}
|
||||
if (/\b(?:patente(?:\s+di\s+guida)?|driving\s+licen[cs]e)(?:\s+(?:numero|number|n\.?))?\s*[:#-]?\s*[A-Z0-9]{8,12}\b/iu.test(value)) {
|
||||
return { kind: "content", ruleId: "pii.drivers_license_number" };
|
||||
}
|
||||
for (const match of value.matchAll(/(?<![A-Z0-9])[A-Z]{2}\d{2}(?:\s?[A-Z0-9]){11,30}(?![A-Z0-9])/giu)) {
|
||||
if (validIban(match[0])) return { kind: "content", ruleId: "financial.iban" };
|
||||
}
|
||||
for (const match of value.matchAll(/(?<!\d)(?:\d[ -]?){13,19}(?!\d)/gu)) {
|
||||
if (validPaymentCard(match[0])) {
|
||||
return { kind: "content", ruleId: "financial.payment_card" };
|
||||
}
|
||||
}
|
||||
for (const match of value.matchAll(/(?<![A-Z0-9])[A-Z]{6}[A-Z0-9]{2}(?:[A-Z0-9]{3})?(?![A-Z0-9])/giu)) {
|
||||
const before = value.slice(Math.max(0, (match.index ?? 0) - 24), match.index ?? 0);
|
||||
if (/\b(?:bic|swift)\s*[:=-]?\s*$/iu.test(before) && validator.isBIC(match[0])) {
|
||||
return { kind: "content", ruleId: "financial.bic" };
|
||||
}
|
||||
}
|
||||
for (const match of value.matchAll(/(?<!\d)(?:IT[ .-]?)?\d{11}(?!\d)/giu)) {
|
||||
const candidate = match[0].replace(/[ .-]/gu, "");
|
||||
if (validator.isVAT(candidate.replace(/^IT/iu, ""), "IT")) {
|
||||
return { kind: "content", ruleId: "pii.italian_vat" };
|
||||
}
|
||||
}
|
||||
for (const match of value.matchAll(/(?<![A-F0-9])(?:[A-F0-9]{2}[:-]){5}[A-F0-9]{2}(?![A-F0-9])/giu)) {
|
||||
if (validator.isMACAddress(match[0])) {
|
||||
return { kind: "content", ruleId: "network.mac_address" };
|
||||
}
|
||||
}
|
||||
for (const match of value.matchAll(/(?<![A-F0-9:.])[A-F0-9:.]{3,45}(?![A-F0-9:.])/giu)) {
|
||||
if (validator.isIP(match[0])) return { kind: "content", ruleId: "network.ip_address" };
|
||||
}
|
||||
for (const match of value.matchAll(/(?<![A-F0-9-])[0-9A-F]{8}-[0-9A-F]{4}-[1-8][0-9A-F]{3}-[89AB][0-9A-F]{3}-[0-9A-F]{12}(?![A-F0-9-])/giu)) {
|
||||
if (validator.isUUID(match[0])) return { kind: "content", ruleId: "pii.uuid" };
|
||||
}
|
||||
for (const match of value.matchAll(/\b(?:https?|ftp):\/\/[^\s<>"']+/giu)) {
|
||||
const candidate = match[0].replace(/[.,;:!?\])}]+$/u, "");
|
||||
if (validator.isURL(candidate, { require_protocol: true })) {
|
||||
return { kind: "content", ruleId: "network.url" };
|
||||
}
|
||||
}
|
||||
if (findPhoneNumbersInText(value, "IT").some((match) => match.number.isValid())) {
|
||||
return { kind: "content", ruleId: "pii.phone_number" };
|
||||
}
|
||||
if (jsonHasSensitiveKey(value)) {
|
||||
return { kind: "content", ruleId: "pii.json_sensitive_key" };
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
/** Sole decision module for local column-level sensitivity assessments. */
|
||||
export class SensitivityClassifier {
|
||||
constructor(
|
||||
private readonly values: SensitivityValueSource,
|
||||
private readonly detector?: LocalNerDetector,
|
||||
private readonly options: {
|
||||
fullScanBudgetMs?: number;
|
||||
runBudgetMs?: number;
|
||||
nerConfidenceThreshold?: number;
|
||||
maxNerValuesPerColumn?: number;
|
||||
maxNerCandidatesPerTable?: number;
|
||||
now?: () => number;
|
||||
} = {},
|
||||
) {}
|
||||
|
||||
async assessTable(
|
||||
target: SensitivityTableTarget,
|
||||
signal: AbortSignal,
|
||||
runDeadline?: number,
|
||||
nerBudget?: SensitivityNerBudget,
|
||||
): Promise<readonly SensitivityColumnAssessment[]> {
|
||||
const now = this.options.now ?? Date.now;
|
||||
const deadline = runDeadline ?? now() + (this.options.runBudgetMs ?? 60_000);
|
||||
const evidence = new Map(target.columns.map((column) => {
|
||||
const match = metadataEvidence(column);
|
||||
return [column.id, match ? [match] : [] as SensitivityEvidence[]];
|
||||
}));
|
||||
const observed = new Map(target.columns.map((column) => [column.id, 0]));
|
||||
const nerCandidates = new Map(target.columns.map((column) => [column.id, [] as string[]]));
|
||||
const maxNerValuesPerColumn = boundedCount(this.options.maxNerValuesPerColumn, 8, 8);
|
||||
const unsupported = new Set(target.columns
|
||||
.filter((column) => UNSUPPORTED_BINARY_TYPE.test(column.dataType))
|
||||
.map((column) => column.id));
|
||||
const scannableColumns = target.columns.filter((column) => (
|
||||
!unsupported.has(column.id) && evidence.get(column.id)!.length === 0
|
||||
));
|
||||
let coverage: SensitivityScanCoverage = { kind: "unavailable", observedRows: 0 };
|
||||
if (scannableColumns.length > 0 && now() < deadline) {
|
||||
try {
|
||||
coverage = await this.values.scanTable({
|
||||
...target,
|
||||
columns: scannableColumns,
|
||||
fullScanBudgetMs: this.options.fullScanBudgetMs ?? 5_000,
|
||||
deadline,
|
||||
}, (batch) => {
|
||||
for (const item of batch) {
|
||||
if (!evidence.has(item.columnId) || item.value === null) continue;
|
||||
observed.set(item.columnId, (observed.get(item.columnId) ?? 0) + 1);
|
||||
const matches = evidence.get(item.columnId)!;
|
||||
if (matches.length === 0 && (item.characterLength ?? item.value.length) > 500) {
|
||||
matches.push({ kind: "length", ruleId: "text.over_500_characters" });
|
||||
} else if (matches.length === 0) {
|
||||
const match = contentEvidence(item.value);
|
||||
if (match) matches.push(match);
|
||||
else {
|
||||
const candidates = nerCandidates.get(item.columnId)!;
|
||||
if (candidates.length < maxNerValuesPerColumn && !candidates.includes(item.value)) {
|
||||
candidates.push(item.value);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}, signal);
|
||||
} catch (error) {
|
||||
if (!(error instanceof CatalogConnectorError)) throw error;
|
||||
}
|
||||
}
|
||||
|
||||
if (this.detector && (this.detector.isReady?.() ?? true) && !signal.aborted
|
||||
&& now() < deadline && (nerBudget?.remainingMs ?? 1) > 0) {
|
||||
const candidates: LocalNerCandidate[] = [];
|
||||
const maxCandidates = boundedCount(this.options.maxNerCandidatesPerTable, 2, 1_024);
|
||||
candidateSelection: for (let valueIndex = 0; valueIndex < maxNerValuesPerColumn; valueIndex += 1) {
|
||||
for (const column of target.columns) {
|
||||
if (evidence.get(column.id)!.length > 0) continue;
|
||||
const text = nerCandidates.get(column.id)![valueIndex];
|
||||
if (text === undefined) continue;
|
||||
candidates.push({ columnId: column.id, text });
|
||||
if (candidates.length >= maxCandidates) break candidateSelection;
|
||||
}
|
||||
}
|
||||
if (candidates.length > 0) {
|
||||
const threshold = this.options.nerConfidenceThreshold ?? 0.8;
|
||||
const nerStartedAt = now();
|
||||
const allowedNerMs = nerBudget
|
||||
? Math.max(0, nerBudget.remainingMs)
|
||||
: Math.max(0, deadline - nerStartedAt);
|
||||
const nerDeadline = Math.min(deadline, nerStartedAt + allowedNerMs);
|
||||
try {
|
||||
for (let offset = 0; offset < candidates.length; offset += MAX_NER_CANDIDATES_PER_REQUEST) {
|
||||
if (signal.aborted || now() >= nerDeadline) break;
|
||||
try {
|
||||
const detected = await this.detector.detect(
|
||||
candidates.slice(offset, offset + MAX_NER_CANDIDATES_PER_REQUEST),
|
||||
signal,
|
||||
nerDeadline,
|
||||
);
|
||||
for (const item of detected) {
|
||||
const matches = evidence.get(item.columnId);
|
||||
if (!matches || matches.length > 0 || !Number.isFinite(item.confidence)
|
||||
|| item.confidence < threshold || item.confidence > 1) continue;
|
||||
const label = normalizedName(item.label).slice(0, 80);
|
||||
if (!label) continue;
|
||||
matches.push({
|
||||
kind: "ner",
|
||||
ruleId: "ner.entity",
|
||||
label,
|
||||
confidence: item.confidence,
|
||||
});
|
||||
}
|
||||
} catch {
|
||||
// NER is optional: deterministic findings and scan coverage remain authoritative.
|
||||
break;
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
if (nerBudget) {
|
||||
const elapsedMs = Math.max(1, now() - nerStartedAt);
|
||||
nerBudget.remainingMs = Math.max(0, nerBudget.remainingMs - elapsedMs);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return target.columns.map((column) => {
|
||||
const matches = evidence.get(column.id)!;
|
||||
const count = observed.get(column.id) ?? 0;
|
||||
const assessment: SensitivityAssessment = matches.length > 0
|
||||
? "sensitive"
|
||||
: unsupported.has(column.id) || count === 0 || coverage.kind !== "complete"
|
||||
? "unknown"
|
||||
: "non_sensitive";
|
||||
return {
|
||||
columnId: column.id,
|
||||
assessment,
|
||||
proposedSensitive: assessment === "unknown" ? column.sensitive : assessment === "sensitive",
|
||||
evidence: matches.length > 0
|
||||
? matches
|
||||
: assessment === "unknown"
|
||||
? [{
|
||||
kind: "coverage",
|
||||
ruleId: unsupported.has(column.id)
|
||||
? "coverage.unsupported_type"
|
||||
: coverage.kind === "unavailable"
|
||||
? "coverage.unavailable"
|
||||
: count === 0
|
||||
? "coverage.no_values"
|
||||
: "coverage.incomplete",
|
||||
}]
|
||||
: [],
|
||||
observedValues: count,
|
||||
};
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,94 @@
|
||||
import { dirname } from "node:path";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { loadConfig } from "../config.js";
|
||||
import { WorkspaceSecretStore } from "../workspaces/secret-store.js";
|
||||
import { PythonLocalNerDetector } from "./local-ner-detector.js";
|
||||
import { ConcreteCatalogPostgresAccess } from "./postgres-access.js";
|
||||
import { createCatalogRepository } from "./repository.js";
|
||||
import { SensitivityAnalysisService } from "./sensitivity-analysis-service.js";
|
||||
import { SensitivityClassifier } from "./sensitivity-classifier.js";
|
||||
import { ConcreteSensitivityValueSource } from "./sensitivity-value-source.js";
|
||||
|
||||
const WORKSPACE_ID = /^[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?$/u;
|
||||
|
||||
async function main(): Promise<void> {
|
||||
const workspaceId = process.argv[2];
|
||||
if (!workspaceId || !WORKSPACE_ID.test(workspaceId)) {
|
||||
process.stderr.write("Usage: sensitivity-shadow <workspace-id>\n");
|
||||
process.exitCode = 2;
|
||||
return;
|
||||
}
|
||||
|
||||
let detector: PythonLocalNerDetector | undefined;
|
||||
let stage = "configuration";
|
||||
try {
|
||||
const config = loadConfig(process.env);
|
||||
stage = "catalog";
|
||||
const repository = createCatalogRepository(config.catalogDatabase);
|
||||
if (!(await repository.available())) throw new Error("catalog unavailable");
|
||||
const database = await repository.getByWorkspace(workspaceId);
|
||||
if (!database) throw new Error("database unavailable");
|
||||
stage = "source";
|
||||
const secretStore = new WorkspaceSecretStore({
|
||||
root: config.workspaceSecretStoreRoot,
|
||||
runtimeRoot: config.workspaceSecretRuntimeRoot,
|
||||
installationId: config.workspaceRegistry.installationId,
|
||||
});
|
||||
const access = new ConcreteCatalogPostgresAccess(secretStore, {
|
||||
connectTimeoutMs: config.workspaceDiagnosticTimeoutMs,
|
||||
});
|
||||
const source = new ConcreteSensitivityValueSource(access, secretStore);
|
||||
if (config.sensitivityNer) {
|
||||
const workerScript = config.sensitivityNer.workerScript
|
||||
?? fileURLToPath(new URL("../../python/sensitivity_ner_worker.py", import.meta.url));
|
||||
detector = new PythonLocalNerDetector({
|
||||
pythonExecutable: config.sensitivityNer.pythonExecutable,
|
||||
workerScript,
|
||||
modelPath: config.sensitivityNer.modelPath,
|
||||
cwd: dirname(workerScript),
|
||||
threads: config.sensitivityNer.threads,
|
||||
});
|
||||
try {
|
||||
await detector.warmup();
|
||||
} catch {
|
||||
await detector.close();
|
||||
detector = undefined;
|
||||
}
|
||||
}
|
||||
const startedAt = Date.now();
|
||||
stage = "analysis";
|
||||
const suggestions = await new SensitivityAnalysisService(
|
||||
repository,
|
||||
new SensitivityClassifier(source, detector),
|
||||
).analyze(database.id, "all", [], AbortSignal.timeout(65_000));
|
||||
const assessments = { sensitive: 0, nonSensitive: 0, unknown: 0 };
|
||||
const rules = new Map<string, number>();
|
||||
for (const suggestion of suggestions) {
|
||||
if (suggestion.assessment === "sensitive") assessments.sensitive += 1;
|
||||
else if (suggestion.assessment === "non_sensitive") assessments.nonSensitive += 1;
|
||||
else assessments.unknown += 1;
|
||||
for (const evidence of suggestion.evidence) {
|
||||
rules.set(evidence.ruleId, (rules.get(evidence.ruleId) ?? 0) + 1);
|
||||
}
|
||||
}
|
||||
process.stdout.write(`${JSON.stringify({
|
||||
ok: true,
|
||||
policyVersion: "sensitivity-v1",
|
||||
nerEnabled: detector !== undefined,
|
||||
total: suggestions.length,
|
||||
assessments,
|
||||
rules: Object.fromEntries([...rules].sort(([left], [right]) => left.localeCompare(right))),
|
||||
elapsedMs: Date.now() - startedAt,
|
||||
})}\n`);
|
||||
} catch {
|
||||
process.stdout.write(`${JSON.stringify({
|
||||
ok: false,
|
||||
code: `sensitivity_shadow_${stage}_failed`,
|
||||
})}\n`);
|
||||
process.exitCode = 1;
|
||||
} finally {
|
||||
await detector?.close();
|
||||
}
|
||||
}
|
||||
|
||||
await main();
|
||||
@@ -0,0 +1,259 @@
|
||||
import { readFile } from "node:fs/promises";
|
||||
import type { WorkspaceSecretStore } from "../workspaces/secret-store.js";
|
||||
import { CATALOG_SECRET_IDS } from "./secrets.js";
|
||||
import type { CatalogPostgresAccess } from "./postgres-access.js";
|
||||
import type {
|
||||
SensitivityScanCoverage,
|
||||
SensitivityScanRequest,
|
||||
SensitivityValueObservation,
|
||||
SensitivityValueSource,
|
||||
} from "./sensitivity-classifier.js";
|
||||
import { CatalogConnectorError } from "./types.js";
|
||||
|
||||
const MAX_VALUE_CHARACTERS = 501;
|
||||
const DEFAULT_BATCH_ROWS = 200;
|
||||
const DEFAULT_SAMPLE_ROWS = 200;
|
||||
|
||||
function quoteIdentifier(identifier: string): string {
|
||||
return `"${identifier.replaceAll('"', '""')}"`;
|
||||
}
|
||||
|
||||
function projections(request: SensitivityScanRequest): string {
|
||||
return request.columns.flatMap((column, index) => {
|
||||
const identifier = quoteIdentifier(column.name);
|
||||
return [
|
||||
`LEFT((${identifier})::text, ${MAX_VALUE_CHARACTERS}) AS "__value_${index}"`,
|
||||
`CASE WHEN ${identifier} IS NULL THEN NULL ELSE char_length((${identifier})::text) END AS "__length_${index}"`,
|
||||
];
|
||||
}).join(", ");
|
||||
}
|
||||
|
||||
function observations(
|
||||
request: SensitivityScanRequest,
|
||||
rows: readonly Record<string, unknown>[],
|
||||
): SensitivityValueObservation[] {
|
||||
return rows.flatMap((row) => request.columns.map((column, index) => {
|
||||
const sourceValue = row[`__value_${index}`];
|
||||
const sourceLength = row[`__length_${index}`];
|
||||
const value = sourceValue === null || sourceValue === undefined ? null : String(sourceValue);
|
||||
const parsedLength = sourceLength === null || sourceLength === undefined
|
||||
? null
|
||||
: Number(sourceLength);
|
||||
return {
|
||||
columnId: column.id,
|
||||
value,
|
||||
characterLength: parsedLength !== null && Number.isSafeInteger(parsedLength) && parsedLength >= 0
|
||||
? parsedLength
|
||||
: value?.length ?? null,
|
||||
};
|
||||
}));
|
||||
}
|
||||
|
||||
function cancelled(error: unknown): boolean {
|
||||
return Boolean(error && typeof error === "object" && "code" in error && error.code === "57014");
|
||||
}
|
||||
|
||||
interface SensitivityValueSourceOptions {
|
||||
now?: () => number;
|
||||
batchRows?: number;
|
||||
sampleRows?: number;
|
||||
}
|
||||
|
||||
/**
|
||||
* PostgreSQL value adapter. It owns bounded read mechanics and emits normalized values, never a
|
||||
* sensitivity decision.
|
||||
*/
|
||||
export class ConcreteSensitivityValueSource implements SensitivityValueSource {
|
||||
private readonly now: () => number;
|
||||
private readonly batchRows: number;
|
||||
private readonly sampleRows: number;
|
||||
|
||||
constructor(
|
||||
private readonly access: CatalogPostgresAccess,
|
||||
private readonly secretStore?: Pick<WorkspaceSecretStore, "materialize">,
|
||||
options: SensitivityValueSourceOptions = {},
|
||||
) {
|
||||
this.now = options.now ?? Date.now;
|
||||
this.batchRows = options.batchRows ?? DEFAULT_BATCH_ROWS;
|
||||
this.sampleRows = options.sampleRows ?? DEFAULT_SAMPLE_ROWS;
|
||||
}
|
||||
|
||||
async scanTable(
|
||||
request: SensitivityScanRequest,
|
||||
consume: (batch: readonly SensitivityValueObservation[]) => void | Promise<void>,
|
||||
signal: AbortSignal,
|
||||
): Promise<SensitivityScanCoverage> {
|
||||
if (request.columns.length === 0) return { kind: "unavailable", observedRows: 0 };
|
||||
if (request.database.binding.transport === "rest_api") {
|
||||
return await this.scanRest(request, consume, signal);
|
||||
}
|
||||
return await this.scanPostgres(request, consume, signal);
|
||||
}
|
||||
|
||||
private async scanPostgres(
|
||||
request: SensitivityScanRequest,
|
||||
consume: (batch: readonly SensitivityValueObservation[]) => void | Promise<void>,
|
||||
signal: AbortSignal,
|
||||
): Promise<SensitivityScanCoverage> {
|
||||
const client = await this.access.connect(request.database, signal);
|
||||
let transactionOpen = false;
|
||||
const startedAt = this.now();
|
||||
const fullDeadline = Math.min(request.deadline, startedAt + request.fullScanBudgetMs);
|
||||
let observedRows = 0;
|
||||
let cursorOpen = false;
|
||||
try {
|
||||
if (signal.aborted || this.now() >= request.deadline) {
|
||||
return { kind: "sampled", observedRows: 0 };
|
||||
}
|
||||
await client.query("BEGIN TRANSACTION READ ONLY", []);
|
||||
transactionOpen = true;
|
||||
await client.query("SELECT set_config('statement_timeout', $1, true)", [
|
||||
`${Math.max(1, Math.floor(fullDeadline - startedAt))}ms`,
|
||||
]);
|
||||
await client.query("SAVEPOINT sensitivity_full_scan", []);
|
||||
const cursor = [
|
||||
"DECLARE sensitivity_full_scan_cursor NO SCROLL CURSOR FOR",
|
||||
`SELECT ${projections(request)}`,
|
||||
`FROM ${quoteIdentifier(request.database.schema)}.${quoteIdentifier(request.table.name)}`,
|
||||
].join(" ");
|
||||
await client.query(cursor, []);
|
||||
cursorOpen = true;
|
||||
while (!signal.aborted && this.now() < fullDeadline) {
|
||||
let rows: Array<Record<string, unknown>>;
|
||||
try {
|
||||
await client.query("SELECT set_config('statement_timeout', $1, true)", [
|
||||
`${Math.max(1, Math.floor(fullDeadline - this.now()))}ms`,
|
||||
]);
|
||||
rows = (await client.query(
|
||||
`FETCH FORWARD ${this.batchRows} FROM sensitivity_full_scan_cursor`,
|
||||
[],
|
||||
)).rows;
|
||||
} catch (error) {
|
||||
if (!cancelled(error)) throw error;
|
||||
await client.query("ROLLBACK TO SAVEPOINT sensitivity_full_scan", []);
|
||||
cursorOpen = false;
|
||||
break;
|
||||
}
|
||||
if (rows.length > 0) {
|
||||
observedRows += rows.length;
|
||||
await consume(observations(request, rows));
|
||||
}
|
||||
if (rows.length < this.batchRows) {
|
||||
return { kind: "complete", observedRows };
|
||||
}
|
||||
}
|
||||
if (signal.aborted || this.now() >= request.deadline) {
|
||||
return { kind: "sampled", observedRows };
|
||||
}
|
||||
if (cursorOpen) await client.query("CLOSE sensitivity_full_scan_cursor", []);
|
||||
await client.query("RELEASE SAVEPOINT sensitivity_full_scan", []);
|
||||
await client.query("SELECT set_config('statement_timeout', $1, true)", [
|
||||
`${Math.max(1, Math.floor(request.deadline - this.now()))}ms`,
|
||||
]);
|
||||
const sampleSql = [
|
||||
`SELECT ${projections(request)}`,
|
||||
`FROM ${quoteIdentifier(request.database.schema)}.${quoteIdentifier(request.table.name)}`,
|
||||
"TABLESAMPLE SYSTEM (1) REPEATABLE (37)",
|
||||
"LIMIT $1",
|
||||
].join(" ");
|
||||
const sampledRows = (await client.query(sampleSql, [this.sampleRows])).rows;
|
||||
observedRows += sampledRows.length;
|
||||
if (sampledRows.length > 0) await consume(observations(request, sampledRows));
|
||||
return { kind: "sampled", observedRows };
|
||||
} catch (error) {
|
||||
if (error instanceof CatalogConnectorError) throw error;
|
||||
throw new CatalogConnectorError("Sensitivity source scan failed");
|
||||
} finally {
|
||||
if (transactionOpen) await client.query("ROLLBACK", []).catch(() => undefined);
|
||||
await client.end().catch(() => undefined);
|
||||
}
|
||||
}
|
||||
|
||||
private async scanRest(
|
||||
request: SensitivityScanRequest,
|
||||
consume: (batch: readonly SensitivityValueObservation[]) => void | Promise<void>,
|
||||
signal: AbortSignal,
|
||||
): Promise<SensitivityScanCoverage> {
|
||||
if (!this.secretStore) throw new CatalogConnectorError("REST sensitivity scanning is not configured");
|
||||
const auth = request.database.binding.restAuth ?? "bearer";
|
||||
const materialized = this.secretStore.materialize(
|
||||
request.database.workspaceId,
|
||||
auth === "none" ? [] : [CATALOG_SECRET_IDS.apiKey],
|
||||
);
|
||||
const startedAt = this.now();
|
||||
const fullDeadline = Math.min(request.deadline, startedAt + request.fullScanBudgetMs);
|
||||
let observedRows = 0;
|
||||
try {
|
||||
const headers: Record<string, string> = { "content-type": "application/json" };
|
||||
if (auth !== "none") {
|
||||
const credentialFile = materialized.files.get(CATALOG_SECRET_IDS.apiKey);
|
||||
if (!credentialFile) throw new CatalogConnectorError("REST API key is not configured");
|
||||
const credential = (await readFile(credentialFile, "utf8")).trim();
|
||||
if (auth === "bearer") headers.authorization = `Bearer ${credential}`;
|
||||
else headers["x-api-key"] = credential;
|
||||
}
|
||||
const baseUrl = request.database.binding.baseUrl?.replace(/\/+$/u, "");
|
||||
if (!baseUrl) throw new CatalogConnectorError("Database binding is incomplete");
|
||||
const runQuery = async (sql: string, deadline: number): Promise<Array<Record<string, unknown>>> => {
|
||||
const response = await fetch(`${baseUrl}/rpc/run_query`, {
|
||||
method: "POST",
|
||||
headers,
|
||||
body: JSON.stringify({ query_text: sql }),
|
||||
signal: AbortSignal.any([
|
||||
signal,
|
||||
AbortSignal.timeout(Math.max(1, Math.floor(deadline - this.now()))),
|
||||
]),
|
||||
});
|
||||
if (!response.ok) throw new CatalogConnectorError("REST sensitivity source scan failed");
|
||||
const body: unknown = await response.json();
|
||||
if (!Array.isArray(body)
|
||||
|| body.some((row) => !row || typeof row !== "object" || Array.isArray(row))) {
|
||||
throw new CatalogConnectorError("REST sensitivity source response is invalid");
|
||||
}
|
||||
return body as Array<Record<string, unknown>>;
|
||||
};
|
||||
|
||||
let offset = 0;
|
||||
const baseSelect = [
|
||||
`SELECT ${projections(request)}`,
|
||||
`FROM ${quoteIdentifier(request.database.schema)}.${quoteIdentifier(request.table.name)}`,
|
||||
].join(" ");
|
||||
while (!signal.aborted) {
|
||||
let rows: Array<Record<string, unknown>>;
|
||||
try {
|
||||
rows = await runQuery(
|
||||
`${baseSelect} LIMIT ${this.batchRows} OFFSET ${offset}`,
|
||||
fullDeadline,
|
||||
);
|
||||
} catch (error) {
|
||||
if (signal.aborted || this.now() < fullDeadline) throw error;
|
||||
break;
|
||||
}
|
||||
observedRows += rows.length;
|
||||
if (rows.length > 0) await consume(observations(request, rows));
|
||||
if (rows.length < this.batchRows) {
|
||||
return { kind: offset === 0 ? "complete" : "sampled", observedRows };
|
||||
}
|
||||
offset += rows.length;
|
||||
if (this.now() >= fullDeadline) break;
|
||||
}
|
||||
if (signal.aborted || this.now() >= request.deadline) {
|
||||
return { kind: "sampled", observedRows };
|
||||
}
|
||||
const sampleSql = [
|
||||
baseSelect,
|
||||
"TABLESAMPLE SYSTEM (1) REPEATABLE (37)",
|
||||
`LIMIT ${this.sampleRows}`,
|
||||
].join(" ");
|
||||
const sampledRows = await runQuery(sampleSql, request.deadline);
|
||||
observedRows += sampledRows.length;
|
||||
if (sampledRows.length > 0) await consume(observations(request, sampledRows));
|
||||
return { kind: "sampled", observedRows };
|
||||
} catch (error) {
|
||||
if (error instanceof CatalogConnectorError) throw error;
|
||||
throw new CatalogConnectorError("REST sensitivity source scan failed");
|
||||
} finally {
|
||||
materialized.release();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -266,18 +266,21 @@ export interface DescriptionGenerationEvent {
|
||||
createdAt: string;
|
||||
}
|
||||
|
||||
export type SensitiveDataSuggestionScope = "all" | "selected_tables" | "selected_columns";
|
||||
export type SensitiveDataSuggestionStatus = "running" | "completed" | "failed" | "interrupted";
|
||||
export type SensitivityAnalysisScope = "all" | "selected_tables" | "selected_columns";
|
||||
export type SensitivityAnalysisStatus = "running" | "completed" | "failed" | "interrupted";
|
||||
|
||||
export interface SensitiveDataSuggestionRun {
|
||||
export interface SensitivityAnalysisRun {
|
||||
id: string;
|
||||
databaseId: string;
|
||||
scope: SensitiveDataSuggestionScope;
|
||||
modelId: string;
|
||||
status: SensitiveDataSuggestionStatus;
|
||||
scope: SensitivityAnalysisScope;
|
||||
engine: "llm" | "local";
|
||||
modelId: string | null;
|
||||
policyVersion: string | null;
|
||||
status: SensitivityAnalysisStatus;
|
||||
total: number;
|
||||
suggestedSensitive: number;
|
||||
suggestedNonSensitive: number;
|
||||
unknown: number;
|
||||
inputTokens: number;
|
||||
cacheReadTokens: number;
|
||||
outputTokens: number;
|
||||
@@ -288,11 +291,12 @@ export interface SensitiveDataSuggestionRun {
|
||||
errorSummary: string | null;
|
||||
}
|
||||
|
||||
export interface SensitiveDataSuggestionRunUpdate {
|
||||
status?: SensitiveDataSuggestionStatus;
|
||||
export interface SensitivityAnalysisRunUpdate {
|
||||
status?: SensitivityAnalysisStatus;
|
||||
total?: number;
|
||||
suggestedSensitive?: number;
|
||||
suggestedNonSensitive?: number;
|
||||
unknown?: number;
|
||||
finishedAt?: string | null;
|
||||
errorSummary?: string | null;
|
||||
inputTokens?: number;
|
||||
@@ -300,7 +304,7 @@ export interface SensitiveDataSuggestionRunUpdate {
|
||||
outputTokens?: number;
|
||||
}
|
||||
|
||||
export interface SensitiveDataSuggestionEvent {
|
||||
export interface SensitivityAnalysisEvent {
|
||||
runId: string;
|
||||
sequence: number;
|
||||
level: "info" | "warning" | "error";
|
||||
@@ -491,29 +495,29 @@ export interface CatalogRepository {
|
||||
runId: string,
|
||||
afterSequence?: number,
|
||||
): Promise<DescriptionGenerationEvent[]>;
|
||||
createSensitiveDataSuggestionRun(
|
||||
createSensitivityAnalysisRun(
|
||||
databaseId: string,
|
||||
scope: SensitiveDataSuggestionScope,
|
||||
modelId: string,
|
||||
): Promise<SensitiveDataSuggestionRun>;
|
||||
getSensitiveDataSuggestionRun(runId: string): Promise<SensitiveDataSuggestionRun | undefined>;
|
||||
listSensitiveDataSuggestionRuns(limit?: number): Promise<SensitiveDataSuggestionRun[]>;
|
||||
interruptActiveSensitiveDataSuggestionRuns(
|
||||
scope: SensitivityAnalysisScope,
|
||||
origin: { engine: "llm"; modelId: string } | { engine: "local"; policyVersion: string },
|
||||
): Promise<SensitivityAnalysisRun>;
|
||||
getSensitivityAnalysisRun(runId: string): Promise<SensitivityAnalysisRun | undefined>;
|
||||
listSensitivityAnalysisRuns(limit?: number): Promise<SensitivityAnalysisRun[]>;
|
||||
interruptActiveSensitivityAnalysisRuns(
|
||||
errorSummary: string,
|
||||
): Promise<SensitiveDataSuggestionRun[]>;
|
||||
updateSensitiveDataSuggestionRun(
|
||||
): Promise<SensitivityAnalysisRun[]>;
|
||||
updateSensitivityAnalysisRun(
|
||||
runId: string,
|
||||
update: SensitiveDataSuggestionRunUpdate,
|
||||
): Promise<SensitiveDataSuggestionRun | undefined>;
|
||||
appendSensitiveDataSuggestionEvent(
|
||||
update: SensitivityAnalysisRunUpdate,
|
||||
): Promise<SensitivityAnalysisRun | undefined>;
|
||||
appendSensitivityAnalysisEvent(
|
||||
runId: string,
|
||||
level: SensitiveDataSuggestionEvent["level"],
|
||||
level: SensitivityAnalysisEvent["level"],
|
||||
message: string,
|
||||
): Promise<SensitiveDataSuggestionEvent>;
|
||||
listSensitiveDataSuggestionEvents(
|
||||
): Promise<SensitivityAnalysisEvent>;
|
||||
listSensitivityAnalysisEvents(
|
||||
runId: string,
|
||||
afterSequence?: number,
|
||||
): Promise<SensitiveDataSuggestionEvent[]>;
|
||||
): Promise<SensitivityAnalysisEvent[]>;
|
||||
listRelationships(databaseId: string): Promise<CatalogPhysicalRelationship[]>;
|
||||
listLogicalRelationships(databaseId: string): Promise<CatalogLogicalRelationship[]>;
|
||||
getLogicalRelationshipContext(databaseId: string): Promise<CatalogLogicalRelationshipContext | undefined>;
|
||||
|
||||
@@ -33,6 +33,12 @@ export interface AppConfig {
|
||||
secretsFile?: string;
|
||||
installationConfigFile?: string;
|
||||
modelCatalogFile?: string;
|
||||
sensitivityNer?: {
|
||||
pythonExecutable: string;
|
||||
modelPath: string;
|
||||
workerScript?: string;
|
||||
threads: number;
|
||||
};
|
||||
piAuthFile?: string;
|
||||
secretFiles: Readonly<Record<string, string | undefined>>;
|
||||
modelApiKeyFile?: string;
|
||||
@@ -354,6 +360,35 @@ export function loadConfig(
|
||||
|| modelCatalogFile.includes("\0")
|
||||
|| !path.isAbsolute(modelCatalogFile)
|
||||
)) throw new Error("runtime model catalog configuration is invalid");
|
||||
const sensitivityNerModelPath = env.THT_SENSITIVITY_NER_MODEL_PATH;
|
||||
const sensitivityNerPython = env.THT_SENSITIVITY_NER_PYTHON;
|
||||
const sensitivityNerWorker = env.THT_SENSITIVITY_NER_WORKER;
|
||||
for (const [value, label] of [
|
||||
[sensitivityNerModelPath, "model path"],
|
||||
[sensitivityNerPython, "Python executable"],
|
||||
[sensitivityNerWorker, "worker path"],
|
||||
] as const) {
|
||||
if (value !== undefined && (
|
||||
value.length === 0 || value.trim() !== value || value.includes("\0") || !path.isAbsolute(value)
|
||||
)) throw new Error(`sensitivity NER ${label} configuration is invalid`);
|
||||
}
|
||||
if (sensitivityNerModelPath === undefined && (
|
||||
sensitivityNerPython !== undefined
|
||||
|| sensitivityNerWorker !== undefined
|
||||
|| env.THT_SENSITIVITY_NER_THREADS !== undefined
|
||||
)) throw new Error("sensitivity NER settings require a model path");
|
||||
const sensitivityNerThreads = Number(env.THT_SENSITIVITY_NER_THREADS ?? 2);
|
||||
if (!Number.isSafeInteger(sensitivityNerThreads) || sensitivityNerThreads < 1 || sensitivityNerThreads > 8) {
|
||||
throw new Error("sensitivity NER thread configuration is invalid");
|
||||
}
|
||||
const sensitivityNer = sensitivityNerModelPath === undefined
|
||||
? undefined
|
||||
: {
|
||||
modelPath: sensitivityNerModelPath,
|
||||
pythonExecutable: sensitivityNerPython ?? "/opt/sensitivity-ner/bin/python",
|
||||
...(sensitivityNerWorker ? { workerScript: sensitivityNerWorker } : {}),
|
||||
threads: sensitivityNerThreads,
|
||||
};
|
||||
const piAuthFile = env.THT_PI_AUTH_FILE;
|
||||
if (piAuthFile !== undefined && (
|
||||
piAuthFile.trim() !== piAuthFile || piAuthFile.length === 0 || piAuthFile.includes("\0")
|
||||
@@ -442,6 +477,7 @@ export function loadConfig(
|
||||
secretsFile,
|
||||
installationConfigFile,
|
||||
modelCatalogFile,
|
||||
sensitivityNer,
|
||||
piAuthFile,
|
||||
secretFiles,
|
||||
modelApiKeyFile,
|
||||
|
||||
@@ -11,38 +11,35 @@ import {
|
||||
type DescriptionGenerationWorker,
|
||||
} from "../catalog/description-generation-worker.js";
|
||||
import { MetadataGenerationModelUnavailableError } from "../catalog/metadata-generation-models.js";
|
||||
import { ModelCompletionProviderError } from "../catalog/model-completer.js";
|
||||
import {
|
||||
SensitiveDataSuggestionDuplicateTargetIdsError,
|
||||
SensitiveDataSuggestionInvalidResponseError,
|
||||
SensitiveDataSuggestionNoEligibleColumnsError,
|
||||
SensitiveDataSuggestionPayloadTooLargeError,
|
||||
SensitiveDataSuggestionTargetNotFoundError,
|
||||
} from "../catalog/sensitive-data-suggester.js";
|
||||
import type { SensitiveDataSuggestionRunner } from "../catalog/sensitive-data-suggestion-runner.js";
|
||||
SensitivityAnalysisDuplicateTargetIdsError,
|
||||
SensitivityAnalysisInterruptedError,
|
||||
SensitivityAnalysisNoEligibleColumnsError,
|
||||
SensitivityAnalysisTargetNotFoundError,
|
||||
} from "../catalog/sensitivity-analysis-service.js";
|
||||
import type { SensitivityAnalysisRunner } from "../catalog/sensitivity-analysis-runner.js";
|
||||
import {
|
||||
CatalogOperationInProgressError,
|
||||
CatalogConnectorError,
|
||||
CatalogUnavailableError,
|
||||
DescriptionGenerationRunActiveError,
|
||||
type CatalogRepository,
|
||||
type DescriptionGenerationEvent,
|
||||
type DescriptionGenerationRun,
|
||||
type SensitiveDataSuggestionEvent,
|
||||
type SensitiveDataSuggestionRun,
|
||||
type SensitivityAnalysisEvent,
|
||||
type SensitivityAnalysisRun,
|
||||
} from "../catalog/types.js";
|
||||
|
||||
const idSchema = z.uuid();
|
||||
const modelIdSchema = z.string().regex(/^[a-z][a-z0-9._-]{0,63}\/[A-Za-z0-9][A-Za-z0-9._:-]{0,255}$/);
|
||||
const selectedTargetIdsSchema = z.array(idSchema).min(1);
|
||||
const suggestionSchema = z.discriminatedUnion("scope", [
|
||||
z.object({ modelId: modelIdSchema, scope: z.literal("all") }).strict(),
|
||||
z.object({ scope: z.literal("all") }).strict(),
|
||||
z.object({
|
||||
modelId: modelIdSchema,
|
||||
scope: z.literal("selected_tables"),
|
||||
targetIds: selectedTargetIdsSchema,
|
||||
}).strict(),
|
||||
z.object({
|
||||
modelId: modelIdSchema,
|
||||
scope: z.literal("selected_columns"),
|
||||
targetIds: selectedTargetIdsSchema,
|
||||
}).strict(),
|
||||
@@ -112,7 +109,7 @@ function publicRun(run: DescriptionGenerationRun) {
|
||||
};
|
||||
}
|
||||
|
||||
function publicSensitiveDataSuggestionEvent(event: SensitiveDataSuggestionEvent) {
|
||||
function publicSensitivityAnalysisEvent(event: SensitivityAnalysisEvent) {
|
||||
return {
|
||||
runId: event.runId,
|
||||
sequence: event.sequence,
|
||||
@@ -122,16 +119,19 @@ function publicSensitiveDataSuggestionEvent(event: SensitiveDataSuggestionEvent)
|
||||
};
|
||||
}
|
||||
|
||||
function publicSensitiveDataSuggestionRun(run: SensitiveDataSuggestionRun) {
|
||||
function publicSensitivityAnalysisRun(run: SensitivityAnalysisRun) {
|
||||
return {
|
||||
id: run.id,
|
||||
databaseId: run.databaseId,
|
||||
scope: run.scope,
|
||||
engine: run.engine,
|
||||
modelId: run.modelId,
|
||||
policyVersion: run.policyVersion,
|
||||
status: run.status,
|
||||
total: run.total,
|
||||
suggestedSensitive: run.suggestedSensitive,
|
||||
suggestedNonSensitive: run.suggestedNonSensitive,
|
||||
unknown: run.unknown,
|
||||
inputTokens: run.inputTokens,
|
||||
cacheReadTokens: run.cacheReadTokens,
|
||||
outputTokens: run.outputTokens,
|
||||
@@ -232,16 +232,10 @@ function safeSuggestionError(reply: FastifyReply, error: unknown) {
|
||||
if (error instanceof CatalogUnavailableError) {
|
||||
return reply.code(503).send({
|
||||
code: "catalog_unavailable",
|
||||
message: "The database catalog is unavailable, so no sensitive-field suggestions were prepared.",
|
||||
message: "The database catalog is unavailable, so no sensitivity assessments were prepared.",
|
||||
});
|
||||
}
|
||||
if (error instanceof MetadataGenerationModelUnavailableError) {
|
||||
return reply.code(409).send({
|
||||
code: "metadata_generation_model_unavailable",
|
||||
message: "The selected metadata-generation model is unavailable.",
|
||||
});
|
||||
}
|
||||
if (error instanceof SensitiveDataSuggestionTargetNotFoundError) {
|
||||
if (error instanceof SensitivityAnalysisTargetNotFoundError) {
|
||||
const code = error.target === "database"
|
||||
? "database_not_found"
|
||||
: error.target === "table"
|
||||
@@ -254,45 +248,65 @@ function safeSuggestionError(reply: FastifyReply, error: unknown) {
|
||||
: "One or more selected Catalog Columns were not found in this database.";
|
||||
return reply.code(404).send({ code, message });
|
||||
}
|
||||
if (error instanceof SensitiveDataSuggestionDuplicateTargetIdsError) {
|
||||
if (error instanceof SensitivityAnalysisDuplicateTargetIdsError) {
|
||||
return reply.code(400).send({
|
||||
code: "sensitive_data_suggestion_target_ids_duplicate",
|
||||
message: "Each selected table or column must appear only once.",
|
||||
});
|
||||
}
|
||||
if (error instanceof SensitiveDataSuggestionNoEligibleColumnsError) {
|
||||
if (error instanceof SensitivityAnalysisNoEligibleColumnsError) {
|
||||
return reply.code(409).send({
|
||||
code: "sensitive_data_suggestion_no_columns",
|
||||
message: "The selected scope contains no Catalog Columns to classify.",
|
||||
message: "The selected scope contains no Catalog Columns to assess.",
|
||||
});
|
||||
}
|
||||
if (error instanceof SensitiveDataSuggestionPayloadTooLargeError) {
|
||||
return reply.code(413).send({
|
||||
code: "sensitive_data_suggestion_payload_too_large",
|
||||
message: "The selected structural metadata cannot be divided into safe LLM requests.",
|
||||
if (error instanceof SensitivityAnalysisInterruptedError) {
|
||||
return reply.code(504).send({
|
||||
code: "sensitivity_analysis_timeout",
|
||||
message: "Sensitivity analysis reached its time limit. No assessments were applied.",
|
||||
});
|
||||
}
|
||||
if (error instanceof SensitiveDataSuggestionInvalidResponseError) {
|
||||
if (error instanceof CatalogConnectorError) {
|
||||
return reply.code(502).send({
|
||||
code: "sensitive_data_suggestion_invalid_response",
|
||||
message: "The LLM returned an incomplete or invalid classification. No suggestions were applied.",
|
||||
});
|
||||
}
|
||||
if (error instanceof ModelCompletionProviderError) {
|
||||
return reply.code(502).send({
|
||||
code: "sensitive_data_suggestion_provider_unavailable",
|
||||
message: "The selected LLM service could not complete the request. No suggestions were applied.",
|
||||
code: "sensitivity_source_unavailable",
|
||||
message: "The source values could not be inspected safely. No assessments were applied.",
|
||||
});
|
||||
}
|
||||
if (error instanceof z.ZodError) {
|
||||
return reply.code(400).send({
|
||||
code: "sensitive_data_suggestion_request_invalid",
|
||||
message: "Choose a database, one or more tables, or one or more columns to classify.",
|
||||
message: "Choose a database, one or more tables, or one or more columns to assess.",
|
||||
});
|
||||
}
|
||||
return reply.code(500).send({
|
||||
code: "sensitive_data_suggestion_failed",
|
||||
message: "Sensitive-field suggestions failed before review. No changes were applied.",
|
||||
message: "Local sensitivity analysis failed before review. No changes were applied.",
|
||||
});
|
||||
}
|
||||
|
||||
function untilAborted<T>(operation: Promise<T>, signal: AbortSignal): Promise<T> {
|
||||
if (signal.aborted) {
|
||||
void operation.catch(() => undefined);
|
||||
return Promise.reject(new SensitivityAnalysisInterruptedError());
|
||||
}
|
||||
return new Promise<T>((resolve, reject) => {
|
||||
const abort = () => reject(new SensitivityAnalysisInterruptedError());
|
||||
signal.addEventListener("abort", abort, { once: true });
|
||||
if (signal.aborted) {
|
||||
void operation.catch(() => undefined);
|
||||
abort();
|
||||
return;
|
||||
}
|
||||
operation.then(
|
||||
(value) => {
|
||||
signal.removeEventListener("abort", abort);
|
||||
resolve(value);
|
||||
},
|
||||
(error: unknown) => {
|
||||
signal.removeEventListener("abort", abort);
|
||||
reject(error);
|
||||
},
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
@@ -300,18 +314,18 @@ function safeSuggestionHistoryError(reply: FastifyReply, error: unknown) {
|
||||
if (error instanceof CatalogUnavailableError) {
|
||||
return reply.code(503).send({
|
||||
code: "catalog_unavailable",
|
||||
message: "Sensitive Data Suggestion history is unavailable because the database catalog is unavailable.",
|
||||
message: "Sensitivity Analysis history is unavailable because the database catalog is unavailable.",
|
||||
});
|
||||
}
|
||||
if (error instanceof z.ZodError) {
|
||||
return reply.code(400).send({
|
||||
code: "sensitive_data_suggestion_history_request_invalid",
|
||||
message: "Sensitive Data Suggestion history parameters are invalid.",
|
||||
message: "Sensitivity Analysis history parameters are invalid.",
|
||||
});
|
||||
}
|
||||
return reply.code(500).send({
|
||||
code: "sensitive_data_suggestion_history_failed",
|
||||
message: "Sensitive Data Suggestion history could not be loaded.",
|
||||
message: "Sensitivity Analysis history could not be loaded.",
|
||||
});
|
||||
}
|
||||
|
||||
@@ -320,7 +334,7 @@ export function catalogDescriptionGenerationRoutes(
|
||||
deps: {
|
||||
repository: CatalogRepository;
|
||||
worker: DescriptionGenerationWorker;
|
||||
sensitiveDataSuggestionRunner: SensitiveDataSuggestionRunner;
|
||||
sensitivityAnalysisRunner: SensitivityAnalysisRunner;
|
||||
},
|
||||
): void {
|
||||
app.post("/catalog/databases/:databaseId/sensitive-data-suggestions", async (request, reply) => {
|
||||
@@ -328,16 +342,16 @@ export function catalogDescriptionGenerationRoutes(
|
||||
try {
|
||||
const databaseId = idSchema.parse((request.params as { databaseId?: unknown }).databaseId);
|
||||
const input = suggestionSchema.parse(request.body);
|
||||
const result = await deps.sensitiveDataSuggestionRunner.run(
|
||||
const signal = AbortSignal.timeout(60_000);
|
||||
const result = await untilAborted(deps.sensitivityAnalysisRunner.run(
|
||||
databaseId,
|
||||
input.modelId,
|
||||
input.scope,
|
||||
"targetIds" in input ? input.targetIds : [],
|
||||
new AbortController().signal,
|
||||
);
|
||||
signal,
|
||||
), signal);
|
||||
return {
|
||||
suggestions: result.suggestions,
|
||||
run: publicSensitiveDataSuggestionRun(result.run),
|
||||
run: publicSensitivityAnalysisRun(result.run),
|
||||
};
|
||||
} catch (error) {
|
||||
return safeSuggestionError(reply, error);
|
||||
@@ -348,8 +362,8 @@ export function catalogDescriptionGenerationRoutes(
|
||||
if (!manage(request, reply)) return reply;
|
||||
try {
|
||||
const { limit } = historyQuerySchema.parse(request.query);
|
||||
return (await deps.repository.listSensitiveDataSuggestionRuns(limit))
|
||||
.map(publicSensitiveDataSuggestionRun);
|
||||
return (await deps.repository.listSensitivityAnalysisRuns(limit))
|
||||
.map(publicSensitivityAnalysisRun);
|
||||
} catch (error) {
|
||||
return safeSuggestionHistoryError(reply, error);
|
||||
}
|
||||
@@ -359,12 +373,12 @@ export function catalogDescriptionGenerationRoutes(
|
||||
if (!manage(request, reply)) return reply;
|
||||
try {
|
||||
const runId = idSchema.parse((request.params as { runId?: unknown }).runId);
|
||||
const run = await deps.repository.getSensitiveDataSuggestionRun(runId);
|
||||
const run = await deps.repository.getSensitivityAnalysisRun(runId);
|
||||
if (!run) return reply.code(404).send({
|
||||
code: "sensitive_data_suggestion_run_not_found",
|
||||
message: "Sensitive Data Suggestion Run was not found.",
|
||||
message: "Sensitivity Analysis Run was not found.",
|
||||
});
|
||||
return publicSensitiveDataSuggestionRun(run);
|
||||
return publicSensitivityAnalysisRun(run);
|
||||
} catch (error) {
|
||||
return safeSuggestionHistoryError(reply, error);
|
||||
}
|
||||
@@ -375,14 +389,14 @@ export function catalogDescriptionGenerationRoutes(
|
||||
try {
|
||||
const runId = idSchema.parse((request.params as { runId?: unknown }).runId);
|
||||
const { after } = eventQuerySchema.parse(request.query);
|
||||
if (!(await deps.repository.getSensitiveDataSuggestionRun(runId))) {
|
||||
if (!(await deps.repository.getSensitivityAnalysisRun(runId))) {
|
||||
return reply.code(404).send({
|
||||
code: "sensitive_data_suggestion_run_not_found",
|
||||
message: "Sensitive Data Suggestion Run was not found.",
|
||||
message: "Sensitivity Analysis Run was not found.",
|
||||
});
|
||||
}
|
||||
return (await deps.repository.listSensitiveDataSuggestionEvents(runId, after))
|
||||
.map(publicSensitiveDataSuggestionEvent);
|
||||
return (await deps.repository.listSensitivityAnalysisEvents(runId, after))
|
||||
.map(publicSensitivityAnalysisEvent);
|
||||
} catch (error) {
|
||||
return safeSuggestionHistoryError(reply, error);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user