feat(hub): capability connection admin API + UI + release v0.0.32 (#8)

Co-authored-by: Hong Jiarong <me@jrhim.com>
Co-committed-by: Hong Jiarong <me@jrhim.com>
This commit is contained in:
2026-07-18 17:57:04 +08:00
committed by 洪佳荣
parent 5e10419fc8
commit ef96f8d33d
8 changed files with 706 additions and 1 deletions
@@ -0,0 +1,131 @@
/**
* ADR-0027: Admin routes for organization-scoped capability connections.
* GET /api/org/:orgSlug/capability-connections — list all
* GET /api/org/:orgSlug/capability-connections/:capId — read one
* PUT /api/org/:orgSlug/capability-connections/:capId — rotate/create
* DELETE /api/org/:orgSlug/capability-connections/:capId — disable
*/
import type { PrismaClient } from "@prisma/client";
import type { FastifyInstance } from "fastify";
import { CapabilityConnectionService } from "../../capability/capabilityConnectionService.js";
import { CapabilityReadinessError, type CapabilityReadinessProbe } from "../../capability/capabilityReadiness.js";
import type { LocalSecretEnvelope } from "../../security/secretEnvelope.js";
import { requireOrgRole, type GuardDeps } from "../auth/guards.js";
import { handleRouteError } from "../errors.js";
export interface CapabilityConnectionRouteConfig {
readonly prisma: PrismaClient;
readonly sessionSecret: string;
readonly secretEnvelope: LocalSecretEnvelope;
readonly readinessProbe?: CapabilityReadinessProbe;
}
export async function registerCapabilityConnectionRoutes(
app: FastifyInstance,
config: CapabilityConnectionRouteConfig,
): Promise<void> {
const guardDeps: GuardDeps = { prisma: config.prisma, sessionSecret: config.sessionSecret };
const connections = new CapabilityConnectionService(
config.prisma,
config.secretEnvelope,
config.readinessProbe,
);
app.get("/api/org/:orgSlug/capability-connections", async (request, reply) => {
try {
const { orgSlug } = request.params as { orgSlug: string };
const auth = await requireOrgRole(request, reply, guardDeps, { orgSlug });
if (auth === null) return;
return { connections: await connections.list(auth.organization.id) };
} catch (error) {
request.log.error({ requestId: request.id, operation: "capability_connection.list" }, "list failed");
return handleRouteError(reply, error);
}
});
app.get("/api/org/:orgSlug/capability-connections/:capabilityId", async (request, reply) => {
try {
const { orgSlug, capabilityId } = request.params as { orgSlug: string; capabilityId: string };
const auth = await requireOrgRole(request, reply, guardDeps, { orgSlug });
if (auth === null) return;
return { connection: await connections.read(auth.organization.id, capabilityId) };
} catch (error) {
request.log.error({ requestId: request.id, operation: "capability_connection.read" }, "read failed");
return handleRouteError(reply, error);
}
});
app.put("/api/org/:orgSlug/capability-connections/:capabilityId", async (request, reply) => {
try {
const { orgSlug, capabilityId } = request.params as { orgSlug: string; capabilityId: string };
const auth = await requireOrgRole(request, reply, guardDeps, { orgSlug });
if (auth === null) return;
const body = parseBody(request.body);
const result = await connections.rotate({
organizationId: auth.organization.id,
capabilityId,
actorUserId: auth.user.id,
...body,
});
request.log.info({
organizationId: auth.organization.id,
capabilityId,
connectionId: result.id,
status: result.status,
secretVersion: result.activeVersion,
}, result.created ? "Capability Connection created" : "Capability Connection rotated");
const { created, ...metadata } = result;
return reply.status(created ? 201 : 200).send(metadata);
} catch (error) {
const facts = error instanceof CapabilityReadinessError
? {
errorCode: error.code,
failureCategory: error.category,
...(error.upstreamStatus !== undefined ? { upstreamStatus: error.upstreamStatus } : {}),
}
: { errorCode: "capability_connection_write_failed" };
request.log.error({ requestId: request.id, operation: "capability_connection.rotate", ...facts }, "rotate failed");
return handleRouteError(reply, error);
}
});
app.delete("/api/org/:orgSlug/capability-connections/:capabilityId", async (request, reply) => {
try {
const { orgSlug, capabilityId } = request.params as { orgSlug: string; capabilityId: string };
const auth = await requireOrgRole(request, reply, guardDeps, { orgSlug });
if (auth === null) return;
const result = await connections.disable({
organizationId: auth.organization.id,
capabilityId,
actorUserId: auth.user.id,
});
request.log.info({
organizationId: auth.organization.id,
capabilityId,
connectionId: result.id,
status: result.status,
}, "Capability Connection disabled");
return reply.send(result);
} catch (error) {
request.log.error({ requestId: request.id, operation: "capability_connection.disable" }, "disable failed");
return handleRouteError(reply, error);
}
});
}
function parseBody(value: unknown): { readonly accessKeyId: string; readonly accessKeySecret: string; readonly endpoint: string } {
if (typeof value !== "object" || value === null || Array.isArray(value)) {
throw new Error("invalid capability credential body");
}
const body = value as Record<string, unknown>;
for (const name of ["accessKeyId", "accessKeySecret", "endpoint"] as const) {
if (typeof body[name] !== "string" || (body[name] as string).trim() === "") {
throw new Error(`${name} is required`);
}
}
return {
accessKeyId: body["accessKeyId"] as string,
accessKeySecret: body["accessKeySecret"] as string,
endpoint: body["endpoint"] as string,
};
}
+11
View File
@@ -22,6 +22,8 @@ import type { LocalSecretEnvelope } from "../../security/secretEnvelope.js";
import type { ProviderReadinessProbe } from "../../connections/providerReadiness.js";
import type { FeishuReadinessProbe } from "../../connections/feishuReadiness.js";
import { registerFeishuApplicationConnectionRoutes } from "./feishuApplicationConnectionRoutes.js";
import { registerCapabilityConnectionRoutes } from "./capabilityConnectionRoutes.js";
import type { CapabilityReadinessProbe } from "../../capability/capabilityReadiness.js";
export interface OrgRouteConfig {
readonly prisma: PrismaClient;
@@ -30,6 +32,7 @@ export interface OrgRouteConfig {
readonly secretEnvelope: LocalSecretEnvelope;
readonly providerReadinessProbe?: ProviderReadinessProbe;
readonly feishuConnectionReadinessProbe?: FeishuReadinessProbe;
readonly capabilityReadinessProbe?: CapabilityReadinessProbe;
}
export async function registerOrgRoutes(app: FastifyInstance, config: OrgRouteConfig): Promise<void> {
@@ -138,4 +141,12 @@ export async function registerOrgRoutes(app: FastifyInstance, config: OrgRouteCo
? { readinessProbe: config.feishuConnectionReadinessProbe }
: {}),
});
await registerCapabilityConnectionRoutes(app, {
prisma: config.prisma,
sessionSecret: config.sessionSecret,
secretEnvelope: config.secretEnvelope,
...(config.capabilityReadinessProbe !== undefined
? { readinessProbe: config.capabilityReadinessProbe }
: {}),
});
}
@@ -0,0 +1,266 @@
/**
* ADR-0027: Organization-scoped capability connection service. Manages the
* lifecycle (rotate / read / disable) of capability credentials stored in
* ADR-0024 encrypted envelopes with purpose="capability".
*
* Mirrors FeishuApplicationConnectionService, but keyed by (organizationId,
* capabilityId) instead of 1:1 — an org may have multiple capabilities.
*/
import { randomUUID } from "node:crypto";
import type { Prisma, PrismaClient } from "@prisma/client";
import { lockActiveOrganization } from "../org/status.js";
import { LocalSecretEnvelope } from "../security/secretEnvelope.js";
import { probeDocmindCredential, type CapabilityReadinessProbe } from "./capabilityReadiness.js";
import type { CapabilitySecretPayload } from "./types.js";
const CAPABILITY_ID_PATTERN = /^[a-z0-9][a-z0-9._-]{0,63}$/;
const KNOWN_CAPABILITY_IDS = new Set(["pdf_to_md_bundle", "audio_video_to_text"]);
export interface CapabilityCredentialInput {
readonly accessKeyId: string;
readonly accessKeySecret: string;
readonly endpoint: string;
}
export interface RotateCapabilityInput extends CapabilityCredentialInput {
readonly organizationId: string;
readonly capabilityId: string;
readonly actorUserId: string;
}
export interface CapabilityConnectionMetadata {
readonly id: string;
readonly capabilityId: string;
readonly status: "DRAFT" | "ACTIVE" | "DISABLED";
readonly activeVersion: number | null;
readonly keyId: string | null;
readonly createdAt: Date;
readonly updatedAt: Date;
}
export interface CapabilityConnectionWriteResult extends CapabilityConnectionMetadata {
readonly created: boolean;
}
export type CapabilitySecretPayloadV1 = CapabilitySecretPayload;
export class CapabilityConnectionService {
constructor(
private readonly prisma: PrismaClient,
private readonly secrets: LocalSecretEnvelope,
private readonly readinessProbe: CapabilityReadinessProbe = probeDocmindCredential,
) {}
async rotate(input: RotateCapabilityInput): Promise<CapabilityConnectionWriteResult> {
if (!CAPABILITY_ID_PATTERN.test(input.capabilityId)) {
throw new Error(`invalid capabilityId: ${input.capabilityId}`);
}
const payload = validateCredential(input);
await this.prisma.$transaction(async (tx) => {
await requireCapabilityAdmin(tx, input);
});
await this.readinessProbe({
endpoint: payload.endpoint,
accessKeyId: payload.accessKeyId,
accessKeySecret: payload.accessKeySecret,
});
return this.prisma.$transaction(async (tx) => {
await requireCapabilityAdmin(tx, input);
const connection = await tx.organizationCapabilityConnection.upsert({
where: {
organizationId_capabilityId: {
organizationId: input.organizationId,
capabilityId: input.capabilityId,
},
},
update: {},
create: {
id: randomUUID(),
organizationId: input.organizationId,
capabilityId: input.capabilityId,
status: "DRAFT",
},
});
await tx.$queryRaw`SELECT "id" FROM "OrganizationCapabilityConnection" WHERE "id" = ${connection.id} FOR UPDATE`;
const locked = await tx.organizationCapabilityConnection.findUniqueOrThrow({
where: { id: connection.id },
include: {
activeSecretVersion: true,
secretVersions: { orderBy: { version: "desc" }, take: 1, select: { version: true } },
},
});
if (locked.organizationId !== input.organizationId) {
throw new Error("Capability Connection scope changed during rotation");
}
const version = (locked.secretVersions[0]?.version ?? 0) + 1;
const secretVersionId = randomUUID();
const envelope = this.secrets.encryptJson(
{
purpose: "capability",
organizationId: input.organizationId,
connectionId: locked.id,
secretVersionId,
},
payload,
);
const now = new Date();
const secretVersion = await tx.capabilityCredentialVersion.create({
data: {
id: secretVersionId,
connectionId: locked.id,
version,
envelopeVersion: envelope.version,
keyId: envelope.keyId,
envelope: envelope as unknown as Prisma.InputJsonValue,
createdByUserId: input.actorUserId,
},
});
if (locked.activeSecretVersion !== null) {
await tx.capabilityCredentialVersion.update({
where: { id: locked.activeSecretVersion.id },
data: { retiredAt: now },
});
}
const activated = await tx.organizationCapabilityConnection.update({
where: { id: locked.id },
data: {
status: "ACTIVE",
activeSecretVersionId: secretVersion.id,
activatedAt: now,
disabledAt: null,
},
});
await tx.auditEntry.create({
data: {
organizationId: input.organizationId,
actorUserId: input.actorUserId,
action: version === 1 ? "capability.created" : "capability.rotated",
metadata: {
connectionId: locked.id,
capabilityId: input.capabilityId,
status: "ACTIVE",
secretVersion: version,
keyId: envelope.keyId,
},
},
});
return {
...toMetadata(activated, { version, keyId: secretVersion.keyId }),
created: version === 1,
};
});
}
async list(organizationId: string): Promise<CapabilityConnectionMetadata[]> {
const connections = await this.prisma.organizationCapabilityConnection.findMany({
where: { organizationId },
include: { activeSecretVersion: { select: { version: true, keyId: true } } },
orderBy: { capabilityId: "asc" },
});
return connections.map((c) => toMetadata(c, c.activeSecretVersion));
}
async read(organizationId: string, capabilityId: string): Promise<CapabilityConnectionMetadata | null> {
const connection = await this.prisma.organizationCapabilityConnection.findFirst({
where: { organizationId, capabilityId },
include: { activeSecretVersion: { select: { version: true, keyId: true } } },
});
return connection === null ? null : toMetadata(connection, connection.activeSecretVersion);
}
async disable(input: {
readonly organizationId: string;
readonly capabilityId: string;
readonly actorUserId: string;
}): Promise<CapabilityConnectionMetadata> {
return this.prisma.$transaction(async (tx) => {
await requireCapabilityAdmin(tx, input);
const connection = await tx.organizationCapabilityConnection.findFirst({
where: { organizationId: input.organizationId, capabilityId: input.capabilityId },
select: { id: true },
});
if (connection === null) throw new Error("Capability Connection not found");
await tx.$queryRaw`SELECT "id" FROM "OrganizationCapabilityConnection" WHERE "id" = ${connection.id} FOR UPDATE`;
const locked = await tx.organizationCapabilityConnection.findUniqueOrThrow({
where: { id: connection.id },
include: { activeSecretVersion: { select: { version: true, keyId: true } } },
});
if (locked.status === "DISABLED") return toMetadata(locked, locked.activeSecretVersion);
const disabled = await tx.organizationCapabilityConnection.update({
where: { id: locked.id },
data: { status: "DISABLED", disabledAt: new Date() },
});
await tx.auditEntry.create({
data: {
organizationId: input.organizationId,
actorUserId: input.actorUserId,
action: "capability.disabled",
metadata: {
connectionId: locked.id,
capabilityId: input.capabilityId,
previousStatus: locked.status,
status: "DISABLED",
},
},
});
return toMetadata(disabled, locked.activeSecretVersion);
});
}
}
function validateCredential(input: RotateCapabilityInput): CapabilitySecretPayloadV1 {
if (!KNOWN_CAPABILITY_IDS.has(input.capabilityId)) {
throw new Error(`unsupported capabilityId: ${input.capabilityId}`);
}
const accessKeyId = nonEmpty(input.accessKeyId, "accessKeyId");
const accessKeySecret = nonEmpty(input.accessKeySecret, "accessKeySecret");
const endpoint = nonEmpty(input.endpoint, "endpoint");
return { schemaVersion: 1, accessKeyId, accessKeySecret, endpoint };
}
function toMetadata(
connection: {
readonly id: string;
readonly capabilityId: string;
readonly status: string;
readonly createdAt: Date;
readonly updatedAt: Date;
},
secret: { readonly version: number; readonly keyId: string } | null,
): CapabilityConnectionMetadata {
return {
id: connection.id,
capabilityId: connection.capabilityId,
status: connection.status as CapabilityConnectionMetadata["status"],
activeVersion: secret?.version ?? null,
keyId: secret?.keyId ?? null,
createdAt: connection.createdAt,
updatedAt: connection.updatedAt,
};
}
async function requireCapabilityAdmin(
tx: Prisma.TransactionClient,
input: { readonly organizationId: string; readonly actorUserId: string },
): Promise<void> {
await lockActiveOrganization(tx, input.organizationId);
const membership = await tx.organizationMembership.findFirst({
where: {
organizationId: input.organizationId,
userId: input.actorUserId,
role: { in: ["OWNER", "ADMIN"] },
revokedAt: null,
},
select: { id: true },
});
if (membership === null) {
throw new Error("only Organization OWNER or ADMIN may manage capability connections");
}
}
function nonEmpty(value: string, label: string): string {
const trimmed = value.trim();
if (trimmed === "") throw new Error(`${label} must not be empty`);
return trimmed;
}
+69
View File
@@ -0,0 +1,69 @@
/**
* ADR-0027: Capability readiness probe. Validates the Alibaba Cloud docmind
* credential by calling QueryDocParserStatus with a dummy id — a 400 (bad
* request) means the credential is valid (the API accepted auth but rejected
* the id); a 401/403 means the credential is bad.
*/
import { classifyNetworkFailure, type NetworkFailureCategory } from "../connections/networkFailure.js";
export interface CapabilityReadinessInput {
readonly endpoint: string;
readonly accessKeyId: string;
readonly accessKeySecret: string;
}
export type CapabilityReadinessProbe = (input: CapabilityReadinessInput) => Promise<void>;
export class CapabilityReadinessError extends Error {
constructor(
readonly code: "capability_readiness_unsupported" | "capability_readiness_unreachable" | "capability_readiness_rejected",
message: string,
readonly category: NetworkFailureCategory | "configuration" | "http",
readonly upstreamStatus?: number,
) {
super(message);
this.name = "CapabilityReadinessError";
}
}
/**
* Probe the Alibaba Cloud docmind credential. We call QueryDocParserStatus
* with a dummy id. The API will return:
* - 400 (InvalidParameter) → credential valid, just a bad id → probe passes
* - 401/403 (InvalidAccessKey/Forbidden) → credential invalid → probe fails
* - network error → unreachable
*/
export const probeDocmindCredential: CapabilityReadinessProbe = async (input) => {
const url = `https://${input.endpoint}/?Action=QueryDocParserStatus&Id=probe-test&Version=2022-07-11`;
const authHeader = makeBasicAuth(input.accessKeyId, input.accessKeySecret);
let response: Response;
try {
response = await fetch(url, {
method: "GET",
headers: { authorization: authHeader, accept: "application/json" },
redirect: "manual",
signal: AbortSignal.timeout(10_000),
});
} catch (error) {
throw new CapabilityReadinessError(
"capability_readiness_unreachable",
"docmind credential readiness check could not reach the API",
classifyNetworkFailure(error),
);
}
await response.body?.cancel();
if (response.status === 401 || response.status === 403) {
throw new CapabilityReadinessError(
"capability_readiness_rejected",
`docmind credential rejected: status ${response.status}`,
"http",
response.status,
);
}
};
function makeBasicAuth(accessKeyId: string, accessKeySecret: string): string {
const credentials = Buffer.from(`${accessKeyId}:${accessKeySecret}`).toString("base64");
return `Basic ${credentials}`;
}