This commit is contained in:
@@ -3,6 +3,18 @@ export interface ConfiguredTransportUrlOptions {
|
||||
originOnly?: boolean;
|
||||
}
|
||||
|
||||
export function parseCredentialFreeHttpUrl(value: string): URL | undefined {
|
||||
let url: URL;
|
||||
try {
|
||||
url = new URL(value);
|
||||
} catch {
|
||||
return undefined;
|
||||
}
|
||||
if (!["http:", "https:"].includes(url.protocol)
|
||||
|| url.username || url.password || url.search || url.hash) return undefined;
|
||||
return url;
|
||||
}
|
||||
|
||||
function canonicalLoopbackAuthority(value: string): boolean {
|
||||
const match = /^http:\/\/([^/?#]+)(?:[/?#]|$)/.exec(value);
|
||||
if (!match) return false;
|
||||
@@ -27,15 +39,8 @@ export function parseConfiguredTransportUrl(
|
||||
value: string,
|
||||
options: ConfiguredTransportUrlOptions,
|
||||
): URL | undefined {
|
||||
let url: URL;
|
||||
try {
|
||||
url = new URL(value);
|
||||
} catch {
|
||||
return undefined;
|
||||
}
|
||||
if (url.username || url.password || url.search || url.hash || (options.originOnly && url.pathname !== "/")) {
|
||||
return undefined;
|
||||
}
|
||||
const url = parseCredentialFreeHttpUrl(value);
|
||||
if (!url || (options.originOnly && url.pathname !== "/")) return undefined;
|
||||
if (url.protocol === "https:") return url;
|
||||
if (options.allowLoopbackHttp && url.protocol === "http:" && canonicalLoopbackAuthority(value)) return url;
|
||||
return undefined;
|
||||
|
||||
@@ -159,7 +159,7 @@ export class CatalogService {
|
||||
for (const id of Object.values(CATALOG_SECRET_IDS)) this.secretStore.forget(workspaceId, id);
|
||||
}
|
||||
|
||||
async test(database: WorkspaceDatabase): Promise<DatabaseTestResult> {
|
||||
async test(database: WorkspaceDatabase): Promise<WorkspaceDatabase | undefined> {
|
||||
return await this.operations.run(database.id, async () => {
|
||||
const testedAt = new Date().toISOString();
|
||||
const controller = new AbortController();
|
||||
@@ -168,6 +168,7 @@ export class CatalogService {
|
||||
? database.binding.restAuth === "none" ? [] : [CATALOG_SECRET_IDS.apiKey]
|
||||
: [];
|
||||
const materialized = this.secretStore.materialize(database.workspaceId, required);
|
||||
let result: DatabaseTestResult;
|
||||
try {
|
||||
if (database.binding.transport !== "rest_api") {
|
||||
const client = await this.postgres.connect(database, controller.signal);
|
||||
@@ -205,13 +206,13 @@ export class CatalogService {
|
||||
},
|
||||
});
|
||||
}
|
||||
return {
|
||||
result = {
|
||||
connectionStatus: "reachable",
|
||||
testedVersion: database.version,
|
||||
lastTestedAt: testedAt,
|
||||
};
|
||||
} catch {
|
||||
return {
|
||||
result = {
|
||||
connectionStatus: "failed",
|
||||
testedVersion: database.version,
|
||||
lastTestedAt: testedAt,
|
||||
@@ -223,6 +224,7 @@ export class CatalogService {
|
||||
controller.abort();
|
||||
materialized.release();
|
||||
}
|
||||
return await this.repository.recordTest(database.id, database.version, result);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import type { FastifyInstance, FastifyReply, FastifyRequest } from "fastify";
|
||||
import { z } from "zod";
|
||||
import { isPrincipalContext, requirePermission } from "../auth/authorization.js";
|
||||
import { parseCredentialFreeHttpUrl } from "../auth/url-policy.js";
|
||||
import { CatalogService, type CatalogSecretName } from "../catalog/service.js";
|
||||
import { WorkspaceRegistryError } from "../workspaces/git-repository.js";
|
||||
import {
|
||||
@@ -26,7 +27,9 @@ const bindingSchema = z.object({
|
||||
host: optionalText,
|
||||
port: port.optional(),
|
||||
username: optionalText,
|
||||
baseUrl: z.url().max(2048).optional(),
|
||||
baseUrl: z.string().max(2048)
|
||||
.refine((value) => parseCredentialFreeHttpUrl(value) !== undefined)
|
||||
.optional(),
|
||||
restPath: z.string().regex(/^\/(?!\/)[^?#\\\u0000-\u001f]*$/).max(512).optional(),
|
||||
restAuth: z.enum(["none", "bearer", "x-api-key"]).optional(),
|
||||
tlsServername: optionalText,
|
||||
@@ -209,7 +212,7 @@ export function catalogDatabaseRoutes(
|
||||
const database = await deps.repository.get(id);
|
||||
if (!database) return reply.code(404).send({ code: "database_not_found", message: "Database configuration was not found." });
|
||||
if (database.version !== version) return reply.code(409).send({ code: "database_stale", message: "Database configuration changed. Reload and try again." });
|
||||
const tested = await deps.repository.recordTest(id, version, await deps.service.test(database));
|
||||
const tested = await deps.service.test(database);
|
||||
if (!tested) return reply.code(409).send({ code: "database_stale", message: "Database configuration changed. Reload and try again." });
|
||||
return { ...tested, configured: true, secrets: deps.service.configuredSecrets(tested.workspaceId) };
|
||||
} catch (error) { return safeError(reply, error); }
|
||||
|
||||
@@ -5,6 +5,8 @@ import { afterEach, expect, test, vi } from "vitest";
|
||||
import { buildApp } from "../src/app.js";
|
||||
import { loadConfig } from "../src/config.js";
|
||||
import { MemoryCatalogRepository } from "../src/catalog/memory-repository.js";
|
||||
import { CatalogOperationCoordinator } from "../src/catalog/operation-coordinator.js";
|
||||
import type { CatalogPostgresAccess } from "../src/catalog/postgres-access.js";
|
||||
import type { ObservedSchemaSnapshot } from "../src/catalog/types.js";
|
||||
import { WorkspaceSecretStore } from "../src/workspaces/secret-store.js";
|
||||
import type { WorkspaceRegistry, WorkspaceRevision } from "../src/workspaces/registry.js";
|
||||
@@ -28,7 +30,13 @@ const workspace: WorkspaceDescriptor = {
|
||||
};
|
||||
const revision: WorkspaceRevision = { id: "psd-clinical", commit: "a".repeat(40), blob: "b".repeat(40), snapshotPath: "/tmp/psd.yaml" };
|
||||
|
||||
function setup(environment: Record<string, string> = {}) {
|
||||
function setup(
|
||||
environment: Record<string, string> = {},
|
||||
catalogDependencies: {
|
||||
catalogOperationCoordinator?: CatalogOperationCoordinator;
|
||||
catalogPostgresAccess?: CatalogPostgresAccess;
|
||||
} = {},
|
||||
) {
|
||||
const secretRoot = mkdtempSync(join(tmpdir(), "catalog-secret-"));
|
||||
const runtimeRoot = mkdtempSync(join(tmpdir(), "catalog-secret-runtime-"));
|
||||
roots.push(secretRoot, runtimeRoot);
|
||||
@@ -49,6 +57,7 @@ function setup(environment: Record<string, string> = {}) {
|
||||
workspaceSecretStore: secretStore,
|
||||
catalogRepository: repository,
|
||||
workspaceDiagnoser: vi.fn(),
|
||||
...catalogDependencies,
|
||||
});
|
||||
return { app, secretStore, repository };
|
||||
}
|
||||
@@ -130,6 +139,29 @@ test("lists orphaned records and takes the REST diagnostic path from workspace Y
|
||||
]));
|
||||
});
|
||||
|
||||
test.each([
|
||||
"https://reader:secret@psd.example/api",
|
||||
"https://psd.example/api?token=secret",
|
||||
"https://psd.example/api#secret",
|
||||
"ftp://psd.example/api",
|
||||
])("rejects unsafe REST base URL %s before persistence", async (baseUrl) => {
|
||||
const { app, repository } = setup();
|
||||
const response = await app.inject({
|
||||
method: "POST",
|
||||
url: "/catalog/databases",
|
||||
payload: {
|
||||
...direct,
|
||||
binding: { transport: "rest_api", baseUrl, restPath: "/health", restAuth: "bearer" },
|
||||
},
|
||||
});
|
||||
expect(response.statusCode).toBe(400);
|
||||
expect(response.json()).toEqual({
|
||||
code: "database_invalid",
|
||||
message: "Database configuration is invalid.",
|
||||
});
|
||||
expect(await repository.getByWorkspace("psd-clinical")).toBeUndefined();
|
||||
});
|
||||
|
||||
test("uses optimistic versions, keeps secrets write-only, and hard-deletes only local configuration", async () => {
|
||||
const { app, secretStore } = setup();
|
||||
const created = (await app.inject({ method: "POST", url: "/catalog/databases", payload: direct })).json();
|
||||
@@ -151,6 +183,44 @@ test("uses optimistic versions, keeps secrets write-only, and hard-deletes only
|
||||
expect((await app.inject({ method: "GET", url: "/catalog/databases" })).json()).toMatchObject([{ configured: false }]);
|
||||
});
|
||||
|
||||
test("rejects a connection test while another catalog operation owns the database", async () => {
|
||||
const coordinator = new CatalogOperationCoordinator();
|
||||
const postgres: CatalogPostgresAccess = {
|
||||
connect: vi.fn(async () => { throw new Error("connection must not start"); }),
|
||||
};
|
||||
const { app } = setup({}, {
|
||||
catalogOperationCoordinator: coordinator,
|
||||
catalogPostgresAccess: postgres,
|
||||
});
|
||||
const created = (await app.inject({
|
||||
method: "POST",
|
||||
url: "/catalog/databases",
|
||||
payload: direct,
|
||||
})).json();
|
||||
const release = coordinator.reserve(created.id);
|
||||
|
||||
try {
|
||||
const response = await app.inject({
|
||||
method: "POST",
|
||||
url: `/catalog/databases/${created.id}/test`,
|
||||
payload: { version: created.version },
|
||||
});
|
||||
|
||||
expect(response.statusCode).toBe(409);
|
||||
expect(response.json()).toEqual({
|
||||
code: "database_operation_in_progress",
|
||||
message: "A database operation is already in progress.",
|
||||
});
|
||||
expect(postgres.connect).not.toHaveBeenCalled();
|
||||
expect((await app.inject({
|
||||
method: "GET",
|
||||
url: `/catalog/databases/${created.id}`,
|
||||
})).json()).toMatchObject({ connectionStatus: "untested" });
|
||||
} finally {
|
||||
release();
|
||||
}
|
||||
});
|
||||
|
||||
test("returns exact global and per-database fleet metrics", async () => {
|
||||
const { app, repository } = setup();
|
||||
const database = await repository.create(direct);
|
||||
|
||||
@@ -0,0 +1,111 @@
|
||||
import { mkdtempSync, rmSync } from "node:fs";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
import { expect, test } from "vitest";
|
||||
import { MemoryCatalogRepository } from "../src/catalog/memory-repository.js";
|
||||
import { CatalogOperationCoordinator } from "../src/catalog/operation-coordinator.js";
|
||||
import type { CatalogPostgresAccess } from "../src/catalog/postgres-access.js";
|
||||
import { CatalogService } from "../src/catalog/service.js";
|
||||
import type { DatabaseTestResult, WorkspaceDatabase } from "../src/catalog/types.js";
|
||||
import type { WorkspaceRegistry } from "../src/workspaces/registry.js";
|
||||
import { WorkspaceSecretStore } from "../src/workspaces/secret-store.js";
|
||||
|
||||
interface Deferred {
|
||||
promise: Promise<void>;
|
||||
resolve: () => void;
|
||||
}
|
||||
|
||||
function deferred(): Deferred {
|
||||
let resolve!: () => void;
|
||||
const promise = new Promise<void>((done) => { resolve = done; });
|
||||
return { promise, resolve };
|
||||
}
|
||||
|
||||
class PausingRecordRepository extends MemoryCatalogRepository {
|
||||
constructor(
|
||||
private readonly recordStarted: Deferred,
|
||||
private readonly continueRecord: Deferred,
|
||||
) {
|
||||
super();
|
||||
}
|
||||
|
||||
override async recordTest(
|
||||
id: string,
|
||||
expectedVersion: number,
|
||||
result: DatabaseTestResult,
|
||||
): Promise<WorkspaceDatabase | undefined> {
|
||||
this.recordStarted.resolve();
|
||||
await this.continueRecord.promise;
|
||||
return await super.recordTest(id, expectedVersion, result);
|
||||
}
|
||||
}
|
||||
|
||||
test("holds the database reservation until the connection result is recorded", async () => {
|
||||
const secretRoot = mkdtempSync(join(tmpdir(), "catalog-service-secret-"));
|
||||
const runtimeRoot = mkdtempSync(join(tmpdir(), "catalog-service-runtime-"));
|
||||
const recordStarted = deferred();
|
||||
const continueRecord = deferred();
|
||||
const repository = new PausingRecordRepository(recordStarted, continueRecord);
|
||||
const coordinator = new CatalogOperationCoordinator();
|
||||
const secretStore = new WorkspaceSecretStore({
|
||||
root: secretRoot,
|
||||
runtimeRoot,
|
||||
installationId: "test",
|
||||
});
|
||||
const postgres: CatalogPostgresAccess = {
|
||||
connect: async () => ({
|
||||
query: async () => ({ rows: [{ database: "warehouse", schema: "datawarehouse" }] }),
|
||||
end: async () => undefined,
|
||||
}),
|
||||
};
|
||||
|
||||
try {
|
||||
const database = await repository.create({
|
||||
workspaceId: "psd-clinical",
|
||||
engine: "postgres",
|
||||
databaseName: "warehouse",
|
||||
schema: "datawarehouse",
|
||||
binding: {
|
||||
transport: "postgres_direct",
|
||||
host: "db.internal",
|
||||
port: 5432,
|
||||
username: "reader",
|
||||
},
|
||||
});
|
||||
const service = new CatalogService(
|
||||
repository,
|
||||
{} as WorkspaceRegistry,
|
||||
secretStore,
|
||||
[],
|
||||
1_000,
|
||||
postgres,
|
||||
coordinator,
|
||||
);
|
||||
|
||||
const connectionTest = service.test(database);
|
||||
const firstCompletedPhase = await Promise.race([
|
||||
recordStarted.promise.then(() => "recording" as const),
|
||||
connectionTest.then(() => "returned" as const),
|
||||
]);
|
||||
expect(firstCompletedPhase).toBe("recording");
|
||||
|
||||
try {
|
||||
await expect(coordinator.run(database.id, async () => "overlapped"))
|
||||
.rejects.toThrow("A database operation is already in progress");
|
||||
} finally {
|
||||
continueRecord.resolve();
|
||||
}
|
||||
|
||||
expect(await connectionTest).toMatchObject({
|
||||
id: database.id,
|
||||
version: 1,
|
||||
connectionStatus: "reachable",
|
||||
testedVersion: 1,
|
||||
});
|
||||
expect(await coordinator.run(database.id, async () => "released")).toBe("released");
|
||||
} finally {
|
||||
continueRecord.resolve();
|
||||
rmSync(secretRoot, { recursive: true, force: true });
|
||||
rmSync(runtimeRoot, { recursive: true, force: true });
|
||||
}
|
||||
});
|
||||
Reference in New Issue
Block a user