460 lines
17 KiB
TypeScript
460 lines
17 KiB
TypeScript
import type { FastifyInstance, FastifyReply, FastifyRequest } from "fastify";
|
|
import { z } from "zod";
|
|
import { isPrincipalContext, requirePermission } from "../auth/authorization.js";
|
|
import {
|
|
DescriptionGenerationDuplicateTargetIdsError,
|
|
DescriptionGenerationNoEligibleTargetsError,
|
|
DescriptionGenerationRunLiveError,
|
|
DescriptionGenerationTargetIdsRequiredError,
|
|
DescriptionGenerationTargetNotFoundError,
|
|
DescriptionGenerationWorkspaceUnavailableError,
|
|
type DescriptionGenerationWorker,
|
|
} from "../catalog/description-generation-worker.js";
|
|
import { MetadataGenerationModelUnavailableError } from "../catalog/metadata-generation-models.js";
|
|
import { ModelCompletionProviderError } from "../catalog/model-completer.js";
|
|
import {
|
|
SensitiveDataSuggester,
|
|
SensitiveDataSuggestionDuplicateTargetIdsError,
|
|
SensitiveDataSuggestionInvalidResponseError,
|
|
SensitiveDataSuggestionNoEligibleColumnsError,
|
|
SensitiveDataSuggestionPayloadTooLargeError,
|
|
SensitiveDataSuggestionTargetNotFoundError,
|
|
} from "../catalog/sensitive-data-suggester.js";
|
|
import {
|
|
CatalogOperationInProgressError,
|
|
CatalogUnavailableError,
|
|
DescriptionGenerationRunActiveError,
|
|
type CatalogRepository,
|
|
type DescriptionGenerationEvent,
|
|
type DescriptionGenerationRun,
|
|
} from "../catalog/types.js";
|
|
|
|
const idSchema = z.uuid();
|
|
const modelIdSchema = z.string().regex(/^[a-z][a-z0-9._-]{0,63}$/);
|
|
const selectedTargetIdsSchema = z.array(idSchema).min(1);
|
|
const suggestionSchema = z.discriminatedUnion("scope", [
|
|
z.object({ modelId: modelIdSchema, 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(),
|
|
]);
|
|
const startSchema = z.discriminatedUnion("scope", [
|
|
z.object({
|
|
modelId: modelIdSchema,
|
|
scope: z.literal("selected_columns"),
|
|
targetIds: selectedTargetIdsSchema,
|
|
}).strict(),
|
|
z.object({
|
|
modelId: modelIdSchema,
|
|
scope: z.literal("selected_tables"),
|
|
targetIds: selectedTargetIdsSchema,
|
|
}).strict(),
|
|
z.object({ modelId: modelIdSchema, scope: z.literal("all") }).strict(),
|
|
z.object({ modelId: modelIdSchema, scope: z.literal("missing") }).strict(),
|
|
]);
|
|
const eventQuerySchema = z.object({
|
|
after: z.coerce.number().int().nonnegative().default(0),
|
|
}).strict();
|
|
const historyQuerySchema = z.object({
|
|
limit: z.coerce.number().int().min(1).max(100).default(50),
|
|
}).strict();
|
|
const terminalStatuses = new Set<DescriptionGenerationRun["status"]>([
|
|
"completed",
|
|
"completed_with_errors",
|
|
"cancelled",
|
|
"failed",
|
|
"interrupted",
|
|
]);
|
|
|
|
function manage(request: FastifyRequest, reply: FastifyReply) {
|
|
return isPrincipalContext(requirePermission(request, reply, "database.manage"));
|
|
}
|
|
|
|
function publicEvent(event: DescriptionGenerationEvent) {
|
|
return {
|
|
sequence: event.sequence,
|
|
level: event.level,
|
|
message: event.message,
|
|
createdAt: event.createdAt,
|
|
};
|
|
}
|
|
|
|
function publicRun(run: DescriptionGenerationRun) {
|
|
return {
|
|
id: run.id,
|
|
databaseId: run.databaseId,
|
|
scope: run.scope,
|
|
modelId: run.modelId,
|
|
language: run.language,
|
|
status: run.status,
|
|
total: run.total,
|
|
processed: run.processed,
|
|
generated: run.generated,
|
|
nonGeneratable: run.nonGeneratable,
|
|
failed: run.failed,
|
|
createdAt: run.createdAt,
|
|
startedAt: run.startedAt,
|
|
updatedAt: run.updatedAt,
|
|
finishedAt: run.finishedAt,
|
|
errorSummary: run.errorSummary,
|
|
};
|
|
}
|
|
|
|
function safeError(reply: FastifyReply, error: unknown) {
|
|
if (error instanceof CatalogUnavailableError) {
|
|
return reply.code(503).send({
|
|
code: "catalog_unavailable",
|
|
message: "Database catalog is unavailable.",
|
|
});
|
|
}
|
|
if (error instanceof DescriptionGenerationRunActiveError) {
|
|
return reply.code(409).send({
|
|
code: "description_generation_run_active",
|
|
message: "A Description Generation Run is already active.",
|
|
});
|
|
}
|
|
if (error instanceof DescriptionGenerationRunLiveError) {
|
|
return reply.code(409).send({
|
|
code: "description_generation_run_live",
|
|
message: "A local Description Generation worker or helper is still running.",
|
|
});
|
|
}
|
|
if (error instanceof CatalogOperationInProgressError) {
|
|
return reply.code(409).send({
|
|
code: "database_operation_in_progress",
|
|
message: "A database operation is already in progress.",
|
|
});
|
|
}
|
|
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 DescriptionGenerationDuplicateTargetIdsError) {
|
|
return reply.code(400).send({
|
|
code: "description_generation_target_ids_duplicate",
|
|
message: "Description generation target IDs must be unique.",
|
|
});
|
|
}
|
|
if (error instanceof DescriptionGenerationTargetIdsRequiredError) {
|
|
return reply.code(400).send({
|
|
code: "description_generation_request_invalid",
|
|
message: "At least one description generation target ID is required.",
|
|
});
|
|
}
|
|
if (error instanceof DescriptionGenerationNoEligibleTargetsError) {
|
|
return reply.code(409).send({
|
|
code: "description_generation_no_eligible_targets",
|
|
message: error.scope === "all"
|
|
? "No Catalog Tables or Catalog Columns are available for description generation."
|
|
: "No Catalog Tables or Catalog Columns have a missing Generated Description.",
|
|
});
|
|
}
|
|
if (error instanceof DescriptionGenerationTargetNotFoundError) {
|
|
const code = error.target === "database"
|
|
? "database_not_found"
|
|
: error.target === "table"
|
|
? "catalog_table_not_found"
|
|
: "catalog_column_not_found";
|
|
const message = error.target === "database"
|
|
? "Database configuration was not found."
|
|
: error.target === "table"
|
|
? "One or more selected Catalog Tables were not found."
|
|
: "One or more selected Catalog Columns were not found.";
|
|
return reply.code(404).send({
|
|
code,
|
|
message,
|
|
});
|
|
}
|
|
if (error instanceof DescriptionGenerationWorkspaceUnavailableError) {
|
|
return reply.code(409).send({
|
|
code: "workspace_configuration_unavailable",
|
|
message: "The database workspace configuration is unavailable.",
|
|
});
|
|
}
|
|
if (error instanceof z.ZodError) {
|
|
return reply.code(400).send({
|
|
code: "description_generation_request_invalid",
|
|
message: "Description generation request is invalid.",
|
|
});
|
|
}
|
|
return reply.code(500).send({
|
|
code: "description_generation_failed",
|
|
message: "Description generation failed.",
|
|
});
|
|
}
|
|
|
|
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.",
|
|
});
|
|
}
|
|
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) {
|
|
const code = error.target === "database"
|
|
? "database_not_found"
|
|
: error.target === "table"
|
|
? "catalog_table_not_found"
|
|
: "catalog_column_not_found";
|
|
const message = error.target === "database"
|
|
? "The database configuration was not found."
|
|
: error.target === "table"
|
|
? "One or more selected Catalog Tables were not found in this database."
|
|
: "One or more selected Catalog Columns were not found in this database.";
|
|
return reply.code(404).send({ code, message });
|
|
}
|
|
if (error instanceof SensitiveDataSuggestionDuplicateTargetIdsError) {
|
|
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) {
|
|
return reply.code(409).send({
|
|
code: "sensitive_data_suggestion_no_columns",
|
|
message: "The selected scope contains no Catalog Columns to classify.",
|
|
});
|
|
}
|
|
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 SensitiveDataSuggestionInvalidResponseError) {
|
|
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.",
|
|
});
|
|
}
|
|
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.",
|
|
});
|
|
}
|
|
return reply.code(500).send({
|
|
code: "sensitive_data_suggestion_failed",
|
|
message: "Sensitive-field suggestions failed before review. No changes were applied.",
|
|
});
|
|
}
|
|
|
|
export function catalogDescriptionGenerationRoutes(
|
|
app: FastifyInstance,
|
|
deps: {
|
|
repository: CatalogRepository;
|
|
worker: DescriptionGenerationWorker;
|
|
sensitiveDataSuggester: SensitiveDataSuggester;
|
|
},
|
|
): void {
|
|
app.post("/catalog/databases/:databaseId/sensitive-data-suggestions", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
try {
|
|
const databaseId = idSchema.parse((request.params as { databaseId?: unknown }).databaseId);
|
|
const input = suggestionSchema.parse(request.body);
|
|
const suggestions = await deps.sensitiveDataSuggester.suggest(
|
|
databaseId,
|
|
input.modelId,
|
|
input.scope,
|
|
"targetIds" in input ? input.targetIds : [],
|
|
new AbortController().signal,
|
|
);
|
|
return { suggestions };
|
|
} catch (error) {
|
|
return safeSuggestionError(reply, error);
|
|
}
|
|
});
|
|
|
|
app.post("/catalog/databases/:databaseId/description-generation-runs", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
try {
|
|
const databaseId = idSchema.parse((request.params as { databaseId?: unknown }).databaseId);
|
|
const input = startSchema.parse(request.body);
|
|
const targetIds = "targetIds" in input ? input.targetIds : [];
|
|
const run = await deps.worker.start(databaseId, input.modelId, input.scope, targetIds);
|
|
return reply.code(202).send(publicRun(run));
|
|
} catch (error) {
|
|
return safeError(reply, error);
|
|
}
|
|
});
|
|
|
|
app.get("/catalog/description-generation-runs", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
try {
|
|
const { limit } = historyQuerySchema.parse(request.query);
|
|
return (await deps.repository.listDescriptionGenerationRuns(limit)).map(publicRun);
|
|
} catch (error) {
|
|
return safeError(reply, error);
|
|
}
|
|
});
|
|
|
|
app.get("/catalog/description-generation-runs/:runId", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
try {
|
|
const runId = idSchema.parse((request.params as { runId?: unknown }).runId);
|
|
const run = await deps.repository.getDescriptionGenerationRun(runId);
|
|
return run ? publicRun(run) : reply.code(404).send({
|
|
code: "description_generation_run_not_found",
|
|
message: "Description Generation Run was not found.",
|
|
});
|
|
} catch (error) {
|
|
return safeError(reply, error);
|
|
}
|
|
});
|
|
|
|
app.post("/catalog/description-generation-runs/:runId/cancel", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
try {
|
|
const runId = idSchema.parse((request.params as { runId?: unknown }).runId);
|
|
const run = await deps.worker.cancel(runId);
|
|
return run ? publicRun(run) : reply.code(404).send({
|
|
code: "description_generation_run_not_found",
|
|
message: "Description Generation Run was not found.",
|
|
});
|
|
} catch (error) {
|
|
return safeError(reply, error);
|
|
}
|
|
});
|
|
|
|
app.post("/catalog/description-generation-runs/unlock", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
try {
|
|
const run = await deps.worker.unlock();
|
|
return run ? publicRun(run) : reply.code(404).send({
|
|
code: "description_generation_run_not_found",
|
|
message: "No stale active Description Generation Run was found.",
|
|
});
|
|
} catch (error) {
|
|
return safeError(reply, error);
|
|
}
|
|
});
|
|
|
|
app.get("/catalog/description-generation-runs/:runId/events", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
let unsubscribe: (() => void) | undefined;
|
|
let hijacked = false;
|
|
try {
|
|
const runId = idSchema.parse((request.params as { runId?: unknown }).runId);
|
|
let after = eventQuerySchema.parse(request.query).after;
|
|
const headerCursor = Number(request.headers["last-event-id"]);
|
|
if (Number.isInteger(headerCursor) && headerCursor >= 0) after = Math.max(after, headerCursor);
|
|
if (!(await deps.repository.getDescriptionGenerationRun(runId))) {
|
|
return reply.code(404).send({
|
|
code: "description_generation_run_not_found",
|
|
message: "Description Generation Run was not found.",
|
|
});
|
|
}
|
|
|
|
const buffered: DescriptionGenerationEvent[] = [];
|
|
let ready = false;
|
|
let closed = false;
|
|
let lastRunSnapshot = "";
|
|
let delivery = Promise.resolve();
|
|
const close = () => {
|
|
if (closed) return;
|
|
closed = true;
|
|
unsubscribe?.();
|
|
if (!reply.raw.destroyed) reply.raw.end();
|
|
};
|
|
const writeRun = (run: DescriptionGenerationRun) => {
|
|
if (closed) return;
|
|
const snapshot = JSON.stringify(publicRun(run));
|
|
if (snapshot === lastRunSnapshot) return;
|
|
lastRunSnapshot = snapshot;
|
|
reply.raw.write(`event: run\ndata: ${snapshot}\n\n`);
|
|
if (terminalStatuses.has(run.status)) close();
|
|
};
|
|
const enqueue = (event: DescriptionGenerationEvent) => {
|
|
delivery = delivery.then(async () => {
|
|
if (closed || event.sequence <= after) return;
|
|
after = event.sequence;
|
|
reply.raw.write(
|
|
`id: ${event.sequence}\nevent: log\ndata: ${JSON.stringify(publicEvent(event))}\n\n`,
|
|
);
|
|
const run = await deps.repository.getDescriptionGenerationRun(runId);
|
|
if (run) writeRun(run);
|
|
}).catch(close);
|
|
};
|
|
unsubscribe = deps.worker.subscribeEvents(runId, (event) => {
|
|
if (ready) enqueue(event);
|
|
else buffered.push(event);
|
|
});
|
|
const persisted = await deps.repository.listDescriptionGenerationEvents(runId, after);
|
|
|
|
reply.hijack();
|
|
hijacked = true;
|
|
reply.raw.writeHead(200, {
|
|
"content-type": "text/event-stream; charset=utf-8",
|
|
"cache-control": "no-cache, no-transform",
|
|
connection: "keep-alive",
|
|
"x-accel-buffering": "no",
|
|
});
|
|
reply.raw.once("close", close);
|
|
request.raw.once("aborted", close);
|
|
for (const event of persisted) {
|
|
if (event.sequence <= after) continue;
|
|
after = event.sequence;
|
|
reply.raw.write(
|
|
`id: ${event.sequence}\nevent: log\ndata: ${JSON.stringify(publicEvent(event))}\n\n`,
|
|
);
|
|
}
|
|
ready = true;
|
|
buffered.sort((a, b) => a.sequence - b.sequence).forEach(enqueue);
|
|
let pending = delivery;
|
|
await pending;
|
|
while (pending !== delivery) {
|
|
pending = delivery;
|
|
await pending;
|
|
}
|
|
const current = await deps.repository.getDescriptionGenerationRun(runId);
|
|
if (current) writeRun(current);
|
|
return reply;
|
|
} catch (error) {
|
|
unsubscribe?.();
|
|
if (hijacked) {
|
|
if (!reply.raw.destroyed) reply.raw.end();
|
|
return reply;
|
|
}
|
|
return safeError(reply, error);
|
|
}
|
|
});
|
|
|
|
app.get("/catalog/description-generation-runs/:runId/events-list", async (request, reply) => {
|
|
if (!manage(request, reply)) return reply;
|
|
try {
|
|
const runId = idSchema.parse((request.params as { runId?: unknown }).runId);
|
|
const { after } = eventQuerySchema.parse(request.query);
|
|
if (!(await deps.repository.getDescriptionGenerationRun(runId))) {
|
|
return reply.code(404).send({
|
|
code: "description_generation_run_not_found",
|
|
message: "Description Generation Run was not found.",
|
|
});
|
|
}
|
|
return (await deps.repository.listDescriptionGenerationEvents(runId, after)).map(publicEvent);
|
|
} catch (error) {
|
|
return safeError(reply, error);
|
|
}
|
|
});
|
|
}
|