Compare commits

...

2 Commits

Author SHA1 Message Date
hongjr03 251ad43ca2 feat(hub): capability connection admin API + UI (ADR-0027)
Adds the admin surface for managing org-scoped capability credentials,
so operators can configure the Aliyun docmind AccessKey from the web UI
instead of needing a CLI or database access.

Backend:
- CapabilityConnectionService: rotate/list/read/disable, mirroring
  FeishuApplicationConnectionService but keyed by (orgId, capabilityId)
- capabilityReadiness: probeDocmindCredential validates the AccessKey by
  calling QueryDocParserStatus with a dummy id (400 = auth ok, 401/403 = bad)
- Admin routes: GET/PUT/DELETE /api/org/:slug/capability-connections/:capId
  + GET list, wired into orgRoutes

Frontend:
- api.ts: CapabilityConnection type + capabilityConnections/rotate/disable
  client methods
- /admin/org/:slug/capabilities page: lists known capabilities
  (pdf_to_md_bundle, audio_video_to_text), shows status/version/keyId,
  inline rotate/disable forms with AccessKey ID/Secret/Endpoint fields
- Nav item '能力' added to sidebar

Version bump 0.0.31 → 0.0.32.
2026-07-18 17:55:54 +08:00
hongjr03 5e10419fc8 fix(hub): docmind client stream upload + correct API response parsing (#7)
Co-authored-by: Hong Jiarong <me@jrhim.com>
Co-committed-by: Hong Jiarong <me@jrhim.com>
2026-07-18 17:26:11 +08:00
10 changed files with 902 additions and 71 deletions
+25
View File
@@ -168,6 +168,16 @@ export interface FeishuApplicationConnection {
updatedAt: string;
}
export interface CapabilityConnection {
id: string;
capabilityId: string;
status: 'DRAFT' | 'ACTIVE' | 'DISABLED';
activeVersion: number | null;
keyId: string | null;
createdAt: string;
updatedAt: string;
}
export interface UsageTotals {
runCount: number;
runsWithCost: number;
@@ -381,6 +391,21 @@ export const api = {
disableFeishuApplication: (slug: string) =>
del(`${orgBase(slug)}/feishu-application-connection`) as Promise<FeishuApplicationConnection>,
capabilityConnections: (slug: string) =>
get(`${orgBase(slug)}/capability-connections`) as Promise<{ connections: CapabilityConnection[] }>,
capabilityConnection: (slug: string, capabilityId: string) =>
get(`${orgBase(slug)}/capability-connections/${encodeURIComponent(capabilityId)}`) as Promise<{
connection: CapabilityConnection | null;
}>,
rotateCapabilityConnection: (
slug: string,
capabilityId: string,
body: { accessKeyId: string; accessKeySecret: string; endpoint: string },
) =>
put(`${orgBase(slug)}/capability-connections/${encodeURIComponent(capabilityId)}`, body) as Promise<CapabilityConnection>,
disableCapabilityConnection: (slug: string, capabilityId: string) =>
del(`${orgBase(slug)}/capability-connections/${encodeURIComponent(capabilityId)}`) as Promise<CapabilityConnection>,
capacityPolicy: (slug: string) => get(`${orgBase(slug)}/capacity-policy`) as Promise<CapacityPolicyView>,
setCapacityPolicy: (slug: string, body: { limits: Partial<Record<CapacityDimension, number | null>> }) =>
put(`${orgBase(slug)}/capacity-policy`, body) as Promise<CapacityPolicyView>,
+1
View File
@@ -28,6 +28,7 @@
{ key: 'skills', label: '技能', icon: 'roles' as const },
{ key: 'roles', label: '角色', icon: 'roles' as const },
{ key: 'feishu', label: '飞书', icon: 'feishu' as const },
{ key: 'capabilities', label: '能力', icon: 'provider' as const },
];
function isAdmin(org: OrgMembership): boolean {
@@ -0,0 +1,202 @@
<script lang="ts">
import { page } from '$app/state';
import { api, type CapabilityConnection } from '$lib/api';
import { fmtDate } from '$lib/format';
import { Label } from 'bits-ui';
import PageHeader from '$lib/components/PageHeader.svelte';
import LoadingState from '$lib/components/LoadingState.svelte';
import ErrorBanner from '$lib/components/ErrorBanner.svelte';
import { toastError, toastSuccess } from '$lib/toast';
const slug = $derived(page.params.slug ?? '');
const KNOWN_CAPABILITIES = [
{ id: 'pdf_to_md_bundle', label: 'PDF → Markdown', description: '将 PDF 转换为带图片的 Markdown bundle(阿里云文档智能,含公式 LaTeX 识别)' },
{ id: 'audio_video_to_text', label: '音视频 → 文本', description: '将音频/视频转写为文本(阿里云文档智能,按秒计费)' },
] as const;
let connections = $state<Map<string, CapabilityConnection>>(new Map());
let loading = $state(true);
let error = $state<string | null>(null);
let editingCap = $state<string | null>(null);
let accessKeyId = $state('');
let accessKeySecret = $state('');
let endpoint = $state('docmind-api.cn-hangzhou.aliyuncs.com');
let saving = $state(false);
let disabling = $state<string | null>(null);
async function load() {
loading = true;
error = null;
try {
const res = await api.capabilityConnections(slug);
connections = new Map(res.connections.map((c) => [c.capabilityId, c]));
} catch (err) {
error = err instanceof Error ? err.message : String(err);
} finally {
loading = false;
}
}
function startEdit(capId: string) {
editingCap = capId;
accessKeyId = '';
accessKeySecret = '';
endpoint = 'docmind-api.cn-hangzhou.aliyuncs.com';
}
function cancelEdit() {
editingCap = null;
}
async function save(capId: string) {
if (accessKeyId.trim() === '' || accessKeySecret.trim() === '' || endpoint.trim() === '') {
toastError('AccessKey ID、AccessKey Secret、Endpoint 均为必填');
return;
}
saving = true;
try {
const result = await api.rotateCapabilityConnection(slug, capId, {
accessKeyId: accessKeyId.trim(),
accessKeySecret: accessKeySecret.trim(),
endpoint: endpoint.trim(),
});
connections.set(capId, result);
connections = new Map(connections);
editingCap = null;
toastSuccess('能力凭据已保存');
} catch (err) {
toastError(err instanceof Error ? err.message : String(err));
} finally {
saving = false;
}
}
async function disable(capId: string) {
if (!confirm('停用后该能力将不可用,确定停用?')) return;
disabling = capId;
try {
const result = await api.disableCapabilityConnection(slug, capId);
connections.set(capId, result);
connections = new Map(connections);
toastSuccess('已停用能力连接');
} catch (err) {
toastError(err instanceof Error ? err.message : String(err));
} finally {
disabling = null;
}
}
function statusBadge(status: string): string {
if (status === 'ACTIVE') return 'saas-badge-primary';
if (status === 'DISABLED') return 'saas-badge-error';
return 'saas-badge-muted';
}
function statusLabel(status: string): string {
if (status === 'ACTIVE') return '已启用';
if (status === 'DISABLED') return '已停用';
return '草稿';
}
$effect(() => {
if (slug) load();
});
</script>
<PageHeader
title="外部能力"
description="管理文档/媒体转换服务的组织级凭据(ADR-0027)。凭据按组织隔离、版本化信封存储,缺失或校验失败即 fail-closed。"
/>
{#if loading}
<LoadingState />
{:else if error}
<ErrorBanner message={error} onretry={load} />
{:else}
<div class="space-y-6">
{#each KNOWN_CAPABILITIES as cap}
{@const conn = connections.get(cap.id)}
<div class="saas-card-pad">
<div class="mb-3 flex items-start justify-between gap-3">
<div>
<div class="flex items-center gap-2">
<h3 class="saas-section-title">{cap.label}</h3>
{#if conn}
<span class={statusBadge(conn.status)}>{statusLabel(conn.status)}</span>
{:else}
<span class="saas-badge-muted">未配置</span>
{/if}
</div>
<p class="saas-muted mt-1 text-sm">{cap.description}</p>
<p class="mt-0.5 font-mono text-xs text-surface-500">{cap.id}</p>
</div>
<div class="flex items-center gap-2">
{#if conn?.status === 'ACTIVE'}
<button
class="saas-btn-ghost text-sm"
onclick={() => disable(cap.id)}
disabled={disabling === cap.id}
>
{disabling === cap.id ? '停用中…' : '停用'}
</button>
{/if}
<button
class="saas-btn-primary text-sm"
onclick={() => startEdit(cap.id)}
disabled={editingCap === cap.id}
>
{conn ? '轮换凭据' : '配置凭据'}
</button>
</div>
</div>
{#if conn}
<dl class="space-y-1.5 text-sm text-surface-700">
<div class="flex justify-between">
<dt class="text-surface-500">版本</dt>
<dd class="font-mono">{conn.activeVersion ?? '—'}</dd>
</div>
<div class="flex justify-between">
<dt class="text-surface-500">密钥 ID</dt>
<dd class="font-mono text-xs">{conn.keyId ?? '—'}</dd>
</div>
<div class="flex justify-between">
<dt class="text-surface-500">更新时间</dt>
<dd>{fmtDate(conn.updatedAt)}</dd>
</div>
</dl>
{/if}
{#if editingCap === cap.id}
<div class="mt-4 border-t border-surface-100 pt-4">
<p class="saas-muted mb-3 text-sm">
阿里云 RAM 用户的 AccessKey。密钥仅写入新版本,旧版本归档。
</p>
<div class="grid gap-4">
<div>
<Label.Root class="saas-label" for="ak-id-{cap.id}">AccessKey ID</Label.Root>
<input id="ak-id-{cap.id}" class="saas-input font-mono text-sm" bind:value={accessKeyId} />
</div>
<div>
<Label.Root class="saas-label" for="ak-secret-{cap.id}">AccessKey Secret</Label.Root>
<input id="ak-secret-{cap.id}" class="saas-input" type="password" bind:value={accessKeySecret} />
</div>
<div>
<Label.Root class="saas-label" for="endpoint-{cap.id}">Endpoint</Label.Root>
<input id="endpoint-{cap.id}" class="saas-input font-mono text-sm" bind:value={endpoint} />
</div>
</div>
<div class="mt-4 flex items-center justify-end gap-3">
<button class="saas-btn-ghost" onclick={cancelEdit} disabled={saving}>取消</button>
<button class="saas-btn-primary" onclick={() => save(cap.id)} disabled={saving}>
{saving ? '保存中…' : '保存'}
</button>
</div>
</div>
{/if}
</div>
{/each}
</div>
{/if}
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@paradigm/hub",
"version": "0.0.31",
"version": "0.0.32",
"private": true,
"type": "module",
"engines": {
+89
View File
@@ -0,0 +1,89 @@
/**
* Smoke test: real Aliyun docmind API call using createReadStream.
* Run: node --experimental-strip-types scripts/smoke-docmind.ts <pdf-path>
*/
import DocmindClient from "@alicloud/docmind-api20220711";
import { RuntimeOptions } from "@alicloud/tea-util";
import { createReadStream } from "node:fs";
import { basename } from "node:path";
const accessKeyId = process.env["ALIBABA_CLOUD_ACCESS_KEY_ID"];
const accessKeySecret = process.env["ALIBABA_CLOUD_ACCESS_KEY_SECRET"];
if (accessKeyId === undefined || accessKeySecret === undefined) {
console.error("Set ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET env vars.");
process.exit(1);
}
const pdfPath = process.argv[2] ?? "../examples/trial/prime_sieve.pdf";
console.log(`Parsing: ${pdfPath}`);
const client = new DocmindClient.default({
endpoint: "docmind-api.cn-hangzhou.aliyuncs.com",
accessKeyId,
accessKeySecret,
type: "access_key",
regionId: "cn-hangzhou",
} as never);
const fileStream = createReadStream(pdfPath);
const runtime = new RuntimeOptions({});
console.log("Submitting job...");
const submitResp = await client.submitDocParserJobAdvance(
new DocmindClient.SubmitDocParserJobAdvanceRequest({
fileUrlObject: fileStream,
fileName: basename(pdfPath),
outputFormat: ["markdown"],
formulaEnhancement: true,
}),
runtime,
);
const jobId = submitResp.body?.data?.id;
console.log("Job ID:", jobId);
if (jobId === undefined) {
console.error("No job id:", JSON.stringify(submitResp.body));
process.exit(1);
}
console.log("Polling...");
const deadline = Date.now() + 5 * 60_000;
let status = "";
let markdownUrl = "";
let pageCount = 0;
while (Date.now() < deadline) {
await new Promise((r) => setTimeout(r, 10_000));
const resp = await client.queryDocParserStatus(
new DocmindClient.QueryDocParserStatusRequest({ id: jobId }),
);
const data = (resp.body as { data?: { status?: string; pageCountEstimate?: number; outputFormatResult?: Array<{ outputFileUrl?: string; outputType?: string }> } }).data;
status = data?.status ?? "";
console.log(` status: ${status}`);
if (status === "success") {
const md = data?.outputFormatResult?.find((r) => r.outputType === "markdown");
markdownUrl = md?.outputFileUrl ?? "";
pageCount = data?.pageCountEstimate ?? 0;
break;
}
if (status === "fail") {
console.error("Job failed!");
process.exit(1);
}
}
if (status !== "success" || markdownUrl === "") {
console.error("Failed or timed out:", status);
process.exit(1);
}
console.log("Downloading markdown from OSS...");
const mdResp = await fetch(markdownUrl);
const markdown = await mdResp.text();
console.log("--- Result ---");
console.log("Pages:", pageCount);
console.log("Cost (USD, est):", ((pageCount > 0 ? pageCount : 1) * 0.0056).toFixed(4));
console.log("Job ID:", jobId);
console.log("Markdown length:", markdown.length, "chars");
console.log("--- Markdown (first 1000 chars) ---");
console.log(markdown.slice(0, 1000));
@@ -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}`;
}
+107 -70
View File
@@ -3,68 +3,49 @@
*
* Uses the official @alicloud/docmind-api20220711 SDK to call the document
* parsing (large model version) API. The API is asynchronous:
* 1. SubmitDocParserJobAdvance — upload local file, get job id
* 2. QueryDocParserStatus — poll until completed
* 3. GetDocParserResult — fetch markdown + page images
* 1. SubmitDocParserJobAdvance — upload local file as a stream, get job id
* 2. QueryDocParserStatus — poll until data.status === "success";
* the markdown output URL is returned in outputFormatResult
* 3. Download markdown from OSS, extract and download referenced images
*
* Pricing (2025-07, aliyun docmind):
* - 图文文档基础链路: 0.02元/页
* - 图文文档增强链路 (含公式 LaTeX): 0.04元/页
* - 视频: 0.002元/秒
* - 音频: 0.00035元/秒
*
* The adapter requests OutputFormat=markdown and FormulaEnhancement=true for
* PDF inputs. Page images are returned as URLs and downloaded into the
* workspace by the adapter.
*/
import $DocmindClient, {
SubmitDocParserJobAdvanceRequest,
QueryDocParserStatusRequest,
GetDocParserResultRequest,
} from "@alicloud/docmind-api20220711";
import { RuntimeOptions } from "@alicloud/tea-util";
import { readFile } from "node:fs/promises";
import { createReadStream } from "node:fs";
import { basename } from "node:path";
import type { CapabilitySecretPayload } from "./types.js";
/** A single extracted image (page render) from the parsed document. */
/** A single extracted image downloaded from the markdown's OSS image URLs. */
export interface DocmindExtractedImage {
/** Suggested relative filename (e.g. "page_1.jpg"). */
readonly filename: string;
/** Raw image bytes (downloaded from the docmind URL). */
readonly data: Uint8Array;
}
/** The structured result of parsing one document. */
export interface DocmindParseResult {
/** Markdown text with image references (relative to output dir). */
readonly markdown: string;
/** Page images extracted from the document, downloaded as bytes. */
readonly images: readonly DocmindExtractedImage[];
/** Number of pages processed (for UsageFact quantity, unit "pages"). */
readonly pageCount: number;
/** USD cost if the service reported one; null if unknown (ADR-0022).
* docmind bills in CNY per page; cost is derived from page count × unit price. */
readonly costUsd: number | null;
/** External job id for the UsageFact correlationId. */
readonly requestId: string | null;
}
/** Options passed to the client. */
export interface DocmindParseOptions {
/** Absolute path to the input file on the Hub's filesystem. */
readonly inputFilePath: string;
}
/**
* Client interface for the Alibaba Cloud Document Mind parsing service.
* The real implementation uses the official SDK; tests inject a mock.
*/
export interface CapabilityProviderClient {
parse(credential: CapabilitySecretPayload, options: DocmindParseOptions): Promise<DocmindParseResult>;
}
/** Errors raised by the docmind client. */
export class DocmindClientError extends Error {
constructor(
message: string,
@@ -76,22 +57,13 @@ export class DocmindClientError extends Error {
}
}
/** Cost per page in USD (0.04 CNY/page ≈ 0.0056 USD, rate ~7.15). Updated when pricing confirmed. */
/** Cost per page in USD (0.04 CNY/page ≈ 0.0056 USD). */
const COST_PER_PAGE_USD = 0.0056;
/** Poll interval for QueryDocParserStatus (ms). Aliyun recommends 10s. */
const POLL_INTERVAL_MS = 10_000;
/** Max poll duration (ms). Aliyun allows 120 minutes; we cap at 5 minutes for a single page bundle. */
const POLL_TIMEOUT_MS = 5 * 60_000;
/** Type alias for the SDK client constructor's config parameter. */
type DocmindConfig = ConstructorParameters<typeof $DocmindClient.default>[0];
/**
* Real Alibaba Cloud Document Mind client using the official SDK.
* Submits a local file, polls for completion, fetches the markdown result,
* and downloads page images.
*/
export class AliyunDocmindClient implements CapabilityProviderClient {
async parse(credential: CapabilitySecretPayload, options: DocmindParseOptions): Promise<DocmindParseResult> {
const config: DocmindConfig = {
@@ -103,11 +75,13 @@ export class AliyunDocmindClient implements CapabilityProviderClient {
} as DocmindConfig;
const client = new $DocmindClient.default(config);
// 1. Submit job with local file upload.
// 1. Submit job with local file as a ReadStream (not a Buffer — the SDK
// serializes Buffers as JSON {type:"Buffer",data:[...]} which the API
// can't read; a Stream is uploaded as multipart form data).
const fileName = basename(options.inputFilePath);
const fileBuffer = await readFile(options.inputFilePath);
const fileStream = createReadStream(options.inputFilePath);
const advanceRequest = new SubmitDocParserJobAdvanceRequest({
fileUrlObject: fileBuffer,
fileUrlObject: fileStream,
fileName,
outputFormat: ["markdown"],
formulaEnhancement: true,
@@ -127,57 +101,112 @@ export class AliyunDocmindClient implements CapabilityProviderClient {
throw new DocmindClientError("docmind submit returned no job id", "docmind_invalid_response");
}
// 2. Poll until completed.
// 2. Poll QueryDocParserStatus until data.status === "success" or "fail".
const deadline = Date.now() + POLL_TIMEOUT_MS;
let completed = false;
let status = "";
let markdownUrl: string | null = null;
let pageCount = 0;
while (Date.now() < deadline) {
await sleep(POLL_INTERVAL_MS);
const statusReq = new QueryDocParserStatusRequest({ id: jobId });
const statusResponse = await client.queryDocParserStatus(statusReq);
const body = statusResponse.body as { completed?: boolean; status?: string };
completed = body.completed === true;
status = body.status ?? "";
if (completed) break;
const body = statusResponse.body as {
data?: {
status?: string;
pageCountEstimate?: number;
outputFormatResult?: Array<{ outputFileUrl?: string; outputType?: string }>;
};
};
status = body.data?.status ?? "";
if (status === "success") {
const mdResult = body.data?.outputFormatResult?.find((r) => r.outputType === "markdown");
markdownUrl = mdResult?.outputFileUrl ?? null;
pageCount = body.data?.pageCountEstimate ?? 0;
break;
}
if (status === "fail") {
throw new DocmindClientError(`docmind job ${jobId} failed`, "docmind_rejected");
}
}
if (!completed) {
if (status !== "success") {
throw new DocmindClientError(`docmind job ${jobId} timed out (status: ${status})`, "docmind_timeout");
}
if (status === "Fail") {
throw new DocmindClientError(`docmind job ${jobId} failed`, "docmind_rejected");
if (markdownUrl === null) {
throw new DocmindClientError("docmind returned no markdown output URL", "docmind_no_output");
}
// 3. Fetch result.
const resultReq = new GetDocParserResultRequest({ id: jobId });
const resultResponse = await client.getDocParserResult(resultReq);
const resultBody = resultResponse.body as {
data?: { markdown?: string; docInfo?: { pages?: Array<{ imageUrl?: string; pageIdCurDoc?: number }> } };
};
const markdown = resultBody.data?.markdown ?? "";
// 3. Download the markdown file from OSS.
let markdown: string;
try {
const mdResp = await fetch(markdownUrl);
if (!mdResp.ok) {
throw new DocmindClientError(`failed to download markdown: HTTP ${mdResp.status}`, "docmind_unreachable");
}
markdown = await mdResp.text();
} catch (e) {
if (e instanceof DocmindClientError) throw e;
throw new DocmindClientError(
e instanceof Error ? e.message : String(e),
"docmind_unreachable",
);
}
if (markdown === "") {
throw new DocmindClientError("docmind returned empty markdown", "docmind_no_output");
}
// 4. Download page images.
const pages = resultBody.data?.docInfo?.pages ?? [];
// 4. Download images referenced in the markdown (OSS URLs).
// Markdown contains ![filename](http://...oss.../image.png?...) entries.
// We rewrite them to local relative paths and download the images.
const images: DocmindExtractedImage[] = [];
for (const page of pages) {
if (page.imageUrl === undefined || page.imageUrl === null || page.imageUrl === "") continue;
const imageRefRegex = /!\[([^\]]*)\]\((https?:\/\/[^)]+)\)/g;
const localMarkdown = markdown.replace(imageRefRegex, (match, altText: string, url: string) => {
const idx = images.findIndex((img) => img.filename === extractFilename(altText, url, images.length));
if (idx >= 0) {
const img = images[idx]!;
return `![${altText}](${img.filename})`;
}
return match;
});
// Collect all image URLs first, then download.
const imageUrls: Array<{ url: string; altText: string }> = [];
let match: RegExpExecArray | null;
const collectRegex = /!\[([^\]]*)\]\((https?:\/\/[^)]+)\)/g;
while ((match = collectRegex.exec(markdown)) !== null) {
imageUrls.push({ url: match[2]!, altText: match[1]! });
}
for (let i = 0; i < imageUrls.length; i++) {
const { url, altText } = imageUrls[i]!;
const filename = extractFilename(altText, url, i);
try {
const resp = await fetch(page.imageUrl);
if (!resp.ok) continue;
const data = new Uint8Array(await resp.arrayBuffer());
const pageNum = page.pageIdCurDoc ?? images.length + 1;
images.push({ filename: `page_${pageNum}.jpg`, data });
const imgResp = await fetch(url);
if (!imgResp.ok) continue;
const data = new Uint8Array(await imgResp.arrayBuffer());
images.push({ filename, data });
} catch {
// Best-effort: skip images that fail to download.
}
}
const pageCount = pages.length > 0 ? pages.length : countMarkdownPages(markdown);
const costUsd = pageCount * COST_PER_PAGE_USD;
// Rewrite markdown with local image paths.
let finalMarkdown = markdown;
let imageIdx = 0;
finalMarkdown = finalMarkdown.replace(imageRefRegex, (match, altText: string, _url: string) => {
if (imageIdx < images.length) {
const img = images[imageIdx]!;
imageIdx++;
return `![${altText}](${img.filename})`;
}
return match;
});
return { markdown, images, pageCount, costUsd, requestId: jobId };
if (pageCount === 0) {
pageCount = Math.max(1, images.length);
}
const costUsd = (pageCount > 0 ? pageCount : 1) * COST_PER_PAGE_USD;
return { markdown: finalMarkdown, images, pageCount, costUsd, requestId: jobId };
}
}
@@ -187,8 +216,16 @@ function sleep(ms: number): Promise<void> {
});
}
/** Fallback page count: count markdown page separators if docInfo is absent. */
function countMarkdownPages(markdown: string): number {
const matches = markdown.match(/\n---\n/g);
return matches !== null ? matches.length + 1 : 1;
/** Derive a clean filename from the alt text or URL. */
function extractFilename(altText: string, url: string, index: number): string {
// Try alt text first (docmind often puts the original filename).
if (altText !== "" && altText.length < 100) {
const cleaned = altText.replace(/[^a-zA-Z0-9._-]/g, "_");
if (cleaned.length > 0) return cleaned;
}
// Fall back to URL path.
const urlPath = new URL(url).pathname;
const base = basename(urlPath);
if (base !== "" && base !== "/") return base;
return `image_${index + 1}.png`;
}