Compare commits

...

5 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
hongjr03 64b3d1fc64 feat(hub): switch capability provider to Aliyun Doc Mind (ADR-0027) (#6)
Co-authored-by: Hong Jiarong <me@jrhim.com>
Co-committed-by: Hong Jiarong <me@jrhim.com>
2026-07-18 16:42:42 +08:00
hongjr03 b673dd1fe9 feat(hub): external capability registry for PDF/ASR transforms (ADR-0027) (#5)
Co-authored-by: Hong Jiarong <me@jrhim.com>
Co-committed-by: Hong Jiarong <me@jrhim.com>
2026-07-18 15:55:02 +08:00
hongjr03 aaa098bb8b feat(hub): usage fact ledger for run-scoped cost attribution (ADR-0026) (#4)
Co-authored-by: Hong Jiarong <me@jrhim.com>
Co-committed-by: Hong Jiarong <me@jrhim.com>
2026-07-18 14:46:41 +08:00
29 changed files with 3129 additions and 48 deletions
+146
View File
@@ -0,0 +1,146 @@
# ADR 0026: Usage Fact Ledger
## Status
Accepted.
## Context
ADR-0022 pins the cost-attribution contract for the platform: token,
provider-reported cost, run count, and duration are attributed by Organization,
Project, Run, model, and Provider Connection; "missing provider cost remains
unknown rather than zero." ADR-0021 scopes usage accounting as operational
reporting, not payment collection (commercial billing stays deferred).
The implementation today stores this as **one scalar per `AgentRun`**:
`inputTokens`, `outputTokens`, `costUsd?`, `costSource?`, written once from the
Claude Agent SDK's `result.total_cost_usd` when the run finishes. This is
adequate for a single provider-reported model loop, but it cannot represent:
- **External capability consumption** inside a run. PDF→Markdown bundle
conversion, audio/video transcription, OCR and similar media transforms are
not the main agent loop. They run as side effects of a run, may use a
different provider/model, may bill in non-token units (pages, seconds,
invocations), and may report cost through a different channel than the
OpenRouter gateway. Today there is no row to write that cost to — it would
either disappear or silently corrupt the run's single scalar.
- **Multiple model calls within one run** (e.g. a sub-model invoked by a
tool, a gateway-side reroute). The scalar collapses them into one number.
- **Pricebook derivation** after the fact. With only a final USD figure and no
`(provider, model, occurredAt, tokens)` fact, an operator cannot re-derive
cost from a price table when the provider did not report it.
Treating each external call as a nested `AgentRun` was considered and
rejected: `AgentRun` carries lock ownership (ADR-0002), session/provider/role
binding (ADR-0017), admission/capacity semantics (ADR-0022), and the
user-visible task boundary. External calls hold none of those. Making them
`AgentRun`s would pollute run counts, admission, lock semantics, and session
continuity, and would still not solve non-token metering.
## Decision
Introduce **`UsageFact`** as the single source of truth for billable
consumption inside an `AgentRun`. An `AgentRun` owns zero or more
append-only `UsageFact` rows; each row records one billable consumption event:
- `kind``model_completion` (the main agent loop) | `external_capability`
| `tool_proxy`. The kind set is `OPEN`; new kinds must be surfaced, not
silently folded into an existing one.
- `provider` — e.g. `openrouter`, `mineru`, `openai_whisper`.
- `model?`, `inputTokens?`, `outputTokens?` — token metering, optional
because non-token capabilities have none.
- `quantity?` + `unit?` — non-token metering (pages, audio_seconds,
invocations), coexisting with tokens rather than replacing them.
- `costUsd?` + `costSource``provider_reported` | `pricebook_derived` |
`unknown`. `costUsd = null` means **unknown, not zero** (ADR-0022). When the
provider reported a cost, `costSource = provider_reported` and that value
wins. When only tokens are known, a later pricebook pass may derive
`costUsd` with `costSource = pricebook_derived`. When neither is possible,
`costSource = unknown` and `costUsd` stays null.
- `occurredAt` — when the consumption happened; the pricebook derivation
depends on this, not on `AgentRun.finishedAt`, because an external
capability may complete before the run finishes.
- `capabilityId?` — for `external_capability` facts, the registered capability
id (e.g. `pdf_to_md_bundle`, `audio_video_to_text`).
- `correlationId?` — external request id for reconciliation / idempotency;
not part of the aggregation key.
### Invariants
1. **Append-only.** A `UsageFact` row is never updated or deleted. A cost
correction is a new row; the old row stays. `onDelete: Cascade` exists only
so a hard run delete (itself not a normal path) cleans up its facts.
2. **Belongs to exactly one Run; never holds a lock.** A fact is a side-effect
ledger of one run, not a sub-run. External capability calls obey the same
boundary: their identity is `capabilityId + correlationId`, not `RunId`.
3. **Missing cost ≠ zero.** Aggregation MUST NOT sum `null` `costUsd` as 0.
A run whose facts all have `costUsd = null` is "cost unknown" — reported as
`runsWithoutCost`, exactly as the pre-migration `costUsd = null` runs are
today. A run with at least one `costUsd`-bearing fact contributes its sum.
### Rollup cache
`AgentRun.costUsd / inputTokens / outputTokens / costSource` columns are kept
as a **derived rollup cache**, not dropped:
- Existing non-`/usage` readers (slash `/usage`, session detail, integration
tests asserting `run.costUsd`) continue to work without code changes for
historical runs.
- On run finish, the writer writes the `UsageFact` row first and then mirrors
it onto `AgentRun` as two separate statements, not one transaction. The
fact is the truth, so it is written first; the cache is derived, so it is
written second. A crash between the two leaves the cache stale, but the
usage service re-reads `UsageFact` directly, so staleness is recoverable
(and the reverse order would lose the truth, which is not). Keeping the
fact insert and the run update in separate statements also avoids holding
the `AgentRun` row lock across the FK ShareLock taken by the insert, which
deadlocks against concurrent workspace teardown under the cascade path
`Organization → Project → AgentRun → UsageFact`.
- The org/project usage service and the slash `/usage` command read from
`UsageFact` directly — they are the canonical aggregation path.
### Migration
A new `UsageFact` table is added. For every existing `AgentRun` with a
non-null `costUsd` or non-null `inputTokens`/`outputTokens`, the migration
inserts **one synthetic `model_completion` fact** carrying the run's
provider, model, tokens, cost, and `costSource`. Its `correlationId` is the
run id, so the row is identifiable as a backfill artefact. This keeps
historical reporting correct under the new aggregation path. Pre-migration
runs with no recorded cost remain `runsWithoutCost` by design, matching
ADR-0022's "missing cost ≠ zero" rule and the existing migration
`20260709143000_agent_run_cost_tracking`'s "no backfill" stance for the
truly unrecorded.
## Consequences
- The billing model now supports external capabilities (PDF→MD, ASR, OCR, …)
without per-capability schema changes: a new capability is one new
`capabilityId` value and one or more `external_capability` facts.
- `AgentRun.costUsd` is no longer the truth; it is a convenience cache.
Readers that need the truth (cost breakdown, per-capability attribution)
must read `UsageFact`. The cache MUST be kept consistent on the write path.
- The capability registry, pricebook, and any commercial settlement remain
`OPEN` and out of pilot scope (ADR-0021). This ADR only pins the ledger
shape and invariants.
- `usage.ts` and `/usage` now perform a join against `UsageFact` rather than a
single-table scan of `AgentRun`; the index on `(runId, occurredAt)` and
`(provider, model, occurredAt)` keeps the existing org/project/report
queries bounded.
- Cost corrections (e.g. a provider rebills a run) produce a new fact row; the
rollup cache must be recomputed. The initial release writes facts once at
run finish and does not support post-hoc correction flows — that remains
`OPEN`.
## Deferred
- **Capability registry** as a first-class spec entity (`Spec.System.Capability`)
with org-scoped enable/disable, input/output contracts, metering schemas,
and org-exclusive credentials (ADR-0024 alignment). This ADR only reserves
the `capabilityId` field and the `external_capability` fact kind.
- **Pricebook** — a versioned price table keyed by `(provider, model, unit)`
with time validity. Required to actually produce `pricebook_derived` costs;
until then, facts without provider-reported cost stay `unknown`.
- **Post-hoc cost correction flow** — append-only today; correction UI and
rollup recompute are `OPEN`.
- **Commercial billing, invoicing, settlement** — still deferred per ADR-0021.
@@ -0,0 +1,152 @@
# ADR 0027: External Capability Registry
## Status
Accepted.
## Context
ADR-0026 introduced the `UsageFact` ledger with a `kind = external_capability`
fact and a `capabilityId` field, but deferred the capability registry itself.
Two concrete needs now force the issue:
- **PDF→Markdown bundle** conversion (and, imminently, audio/video→text)
must run as a side effect of an `AgentRun`, bill in non-token units
(pages, seconds), use a different provider than the model loop, and report
cost through a different channel. It is not a sub-run (ADR-0026 rejected
that) and not a model-provider call (it does not speak the Anthropic/
OpenRouter protocol).
- The Agent already has a Bash tool. Without a first-class capability seam,
the path of least resistance is for the agent to shell out to ad-hoc
scripts that embed API keys, write to arbitrary paths, and report nothing
to the ledger. That is exactly the unattributed, uncontained external
consumption ADR-0022/0026 exist to prevent.
The model-provider connection (`OrganizationProviderConnection`, ADR-0024)
is the wrong seam for these services:
- Its payload schema (`baseUrl` + `authToken` + `anthropicApiKey`) and
readiness probe (`/v1/models?supported_parameters=tools`) are specific to
OpenRouter/Anthropic. MinerU, Whisper, and future OCR/ASR services have
different auth shapes (an API token, optionally a project id) and no
`/v1/models` endpoint.
- Its uniqueness key is `(organizationId, providerId)` where `providerId`
is an OpenRouter-style model-routing id. A capability provider id
(`mineru`) names a *service*, not a model.
- Coupling capability credentials into the model-provider table would force
every capability's auth shape through `ProviderSecretPayloadV1` and every
readiness probe through `probeOpenRouterCredential`.
The Feishu Application Connection (`OrganizationFeishuApplicationConnection`)
is the right structural precedent: it reuses the ADR-0024 envelope (KEK →
DEK → AES-256-GCM, AAD-bound to purpose/org/connection/version) but has its
own connection table, its own payload schema, its own readiness probe, and
its own per-org uniqueness. Capability connections follow the same pattern.
## Decision
### External Capability
An **External Capability** is a platform-registered, org-enabled document or
media transform service invoked as a side effect of an `AgentRun`. It is
identified by a stable `capabilityId` (e.g. `pdf_to_md_bundle`,
`audio_video_to_text`). A capability:
- Has an **input kind** (PDF, image, audio, video, …) and an **output
contract** (markdown bundle with extracted images, transcript text, …).
- Bills in **non-token units** (pages, audio-seconds) recorded on a
`UsageFact` with `kind = external_capability`, or in tokens when the
backing service reports them.
- Writes its output **into the invoking run's workspace** (ADR-0018
`AgentSurface` — no escapes).
- Is invoked through a **capability adapter** in Hub, never by the Agent
shelling out with embedded credentials.
### Capability Connection
Credentials for a capability live in an **`OrganizationCapabilityConnection`**,
structurally identical to the Feishu Application Connection:
- Belongs to exactly one Organization.
- Unique by `(organizationId, capabilityId)`.
- `DRAFT` / `ACTIVE` / `DISABLED`; resolution accepts only `ACTIVE` with a
valid active secret version.
- Secret material is an immutable, AAD-bound, KEK-wrapped envelope version
(`CapabilityCredentialVersion`), reusing the ADR-0024 encryption
machinery with `purpose = "capability"`.
- Its payload schema is capability-specific (`CapabilitySecretPayloadV1`:
`baseUrl`, `apiToken`, optional `projectId`). New capability types extend
the payload, not the connection table.
- A capability-specific **readiness probe** validates the credential before
activation (e.g. MinerU: a trivial authenticated GET). The probe is
injectable, matching the Feishu/provider pattern, so tests never hit the
network.
### Capability Adapter
The adapter is the seam between the Agent and the external service. It:
- Resolves the org's active capability connection (fail-closed, no
process-global fallback — ADR-0024).
- Accepts a workspace-relative input path and an output directory.
- Calls the backing service (MinerU, Whisper, …) via an injectable
`Client` interface so the real HTTP client is swappable and mockable.
- Writes the produced markdown + image assets into the run's workspace.
- Writes one `UsageFact` (or more, if the service reports per-stage
consumption) with `kind = external_capability`, `capabilityId`,
`provider` (the service id), `quantity + unit` (pages / seconds), and
`costUsd + costSource = provider_reported` when the service reports cost.
### Registry
The platform maintains a **registry** of known capabilities: their id,
input kind, output contract, metering unit, and adapter. This is
code-level registration (like `ToolRegistry`), not a database table — a
capability is available to an Organization only when (a) the platform
knows the adapter and (b) the Organization has an `ACTIVE` connection for
it. Both gates are required.
### What is NOT in this ADR
- The capability invocation is **not** a first-class persisted record
(`CapabilityInvocation` table) in this ADR. The `UsageFact` row with
`capabilityId` + `correlationId` is the durable trace. If we later need
a richer invocation log (retries, partial output, multi-stage status),
that is a follow-up; for now the fact is enough.
- **Pricebook** remains deferred (ADR-0026). Capability facts use
`provider_reported` when the service returns cost; otherwise `unknown`.
- **Org-scoped enable/disable policy** beyond connection status is
deferred. An org with an `ACTIVE` connection has the capability; one
without does not. A finer "enabled but no credential" toggle is not
needed yet.
- **Agent-facing tool exposure** (how the Agent discovers and calls the
capability — MCP tool, Bash wrapper, or built-in) is an implementation
detail of the adapter wiring, not a contract concern. The contract pins
that the Agent never receives the capability credential.
## Consequences
- Adding a new external capability (e.g. `image_ocr`) is: register an
adapter, add a `capabilityId` constant, optionally extend the secret
payload — no schema change to `UsageFact` or `AgentRun`.
- The model-provider connection table stays focused on model routing;
capability credentials do not pollute its payload or readiness probe.
- Three connection types now share the ADR-0024 envelope: model-provider,
Feishu application, and capability. Each has its own table, payload
schema, and probe, but the same encryption, rotation, and resolver
boundary.
- The Agent's Bash tool remains available, but the intended path for
document/media transforms is the capability adapter. Whether to narrow
Bash for capability-shaped tasks is an operational policy decision,
not a contract one.
- Tests prove: workspace containment of capability output, fail-closed
credential resolution, `UsageFact` attribution with non-token metering,
and that the Agent process never receives the capability credential.
## Deferred
- `CapabilityInvocation` as a first-class durable record (status, retries,
partial output) — currently the `UsageFact` row is the only trace.
- Pricebook derivation for capability costs (ADR-0026 deferred).
- Org-scoped capability enable/disable policy finer than connection status.
- Agent-facing tool discovery (MCP vs built-in) for capabilities.
+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}
+349 -1
View File
@@ -8,6 +8,9 @@
"name": "@paradigm/hub",
"version": "0.0.31",
"dependencies": {
"@alicloud/credentials": "^2.4.5",
"@alicloud/docmind-api20220711": "^1.4.15",
"@alicloud/tea-util": "^1.4.11",
"@anthropic-ai/claude-agent-sdk": "^0.3.202",
"@fastify/cookie": "^11.0.2",
"@larksuiteoapi/node-sdk": "^1.70.0",
@@ -74,6 +77,175 @@
"zod": "^3.25.76 || ^4.1.8"
}
},
"node_modules/@alicloud/credentials": {
"version": "2.4.5",
"resolved": "https://registry.npmjs.org/@alicloud/credentials/-/credentials-2.4.5.tgz",
"integrity": "sha512-od1ufCxOO7cP2R4EVFliOB0kGo9lUXCibyj/mzmI6yLhxeqhqsegTzVsx5p2NJJsceKJnYcmye7FWKyLJAFBkw==",
"license": "MIT",
"dependencies": {
"@alicloud/tea-typescript": "^1.8.0",
"httpx": "^2.3.3",
"ini": "^1.3.5",
"kitx": "^2.0.0"
}
},
"node_modules/@alicloud/darabonba-array": {
"version": "0.1.2",
"resolved": "https://registry.npmjs.org/@alicloud/darabonba-array/-/darabonba-array-0.1.2.tgz",
"integrity": "sha512-ZPuQ+bJyjrd8XVVm55kl+ypk7OQoi1ZH/DiToaAEQaGvgEjrTcvQkg71//vUX/6cvbLIF5piQDvhrLb+lUEIPQ==",
"license": "ISC",
"dependencies": {
"@alicloud/tea-typescript": "^1.7.1"
}
},
"node_modules/@alicloud/darabonba-encode-util": {
"version": "0.0.2",
"resolved": "https://registry.npmjs.org/@alicloud/darabonba-encode-util/-/darabonba-encode-util-0.0.2.tgz",
"integrity": "sha512-mlsNctkeqmR0RtgE1Rngyeadi5snLOAHBCWEtYf68d7tyKskosXDTNeZ6VCD/UfrUu4N51ItO8zlpfXiOgeg3A==",
"license": "ISC",
"dependencies": {
"moment": "^2.29.1"
}
},
"node_modules/@alicloud/darabonba-map": {
"version": "0.0.1",
"resolved": "https://registry.npmjs.org/@alicloud/darabonba-map/-/darabonba-map-0.0.1.tgz",
"integrity": "sha512-2ep+G3YDvuI+dRYVlmER1LVUQDhf9kEItmVB/bbEu1pgKzelcocCwAc79XZQjTcQGFgjDycf3vH87WLDGLFMlw==",
"license": "ISC",
"dependencies": {
"@alicloud/tea-typescript": "^1.7.1"
}
},
"node_modules/@alicloud/darabonba-signature-util": {
"version": "0.0.4",
"resolved": "https://registry.npmjs.org/@alicloud/darabonba-signature-util/-/darabonba-signature-util-0.0.4.tgz",
"integrity": "sha512-I1TtwtAnzLamgqnAaOkN0IGjwkiti//0a7/auyVThdqiC/3kyafSAn6znysWOmzub4mrzac2WiqblZKFcN5NWg==",
"license": "ISC",
"dependencies": {
"@alicloud/darabonba-encode-util": "^0.0.1"
}
},
"node_modules/@alicloud/darabonba-signature-util/node_modules/@alicloud/darabonba-encode-util": {
"version": "0.0.1",
"resolved": "https://registry.npmjs.org/@alicloud/darabonba-encode-util/-/darabonba-encode-util-0.0.1.tgz",
"integrity": "sha512-Sl5vCRVAYMqwmvXpJLM9hYoCHOMsQlGxaWSGhGWulpKk/NaUBArtoO1B0yHruJf1C5uHhEJIaylYcM48icFHgw==",
"license": "ISC",
"dependencies": {
"@alicloud/tea-typescript": "^1.7.1",
"moment": "^2.29.1"
}
},
"node_modules/@alicloud/darabonba-string": {
"version": "1.0.3",
"resolved": "https://registry.npmjs.org/@alicloud/darabonba-string/-/darabonba-string-1.0.3.tgz",
"integrity": "sha512-NyWwrU8cAIesWk3uHL1Q7pTDTqLkCI/0PmJXC4/4A0MFNAZ9Ouq0iFBsRqvfyUujSSM+WhYLuTfakQXiVLkTMA==",
"license": "Apache-2.0",
"dependencies": {
"@alicloud/tea-typescript": "^1.5.1"
}
},
"node_modules/@alicloud/docmind-api20220711": {
"version": "1.4.15",
"resolved": "https://registry.npmjs.org/@alicloud/docmind-api20220711/-/docmind-api20220711-1.4.15.tgz",
"integrity": "sha512-QmjSDPV52d2B5bd/Dk3Zeofu+cJFPKYS+MphIHqlulfYsOynLgOaDqTRMGTK/RhcmVCwG14CyZ/15GBF00GRFw==",
"license": "Apache-2.0",
"dependencies": {
"@alicloud/credentials": "^2.4.2",
"@alicloud/openapi-core": "^1.0.0",
"@darabonba/typescript": "^1.0.0"
}
},
"node_modules/@alicloud/endpoint-util": {
"version": "0.0.1",
"resolved": "https://registry.npmjs.org/@alicloud/endpoint-util/-/endpoint-util-0.0.1.tgz",
"integrity": "sha512-+pH7/KEXup84cHzIL6UJAaPqETvln4yXlD9JzlrqioyCSaWxbug5FUobsiI6fuUOpw5WwoB3fWAtGbFnJ1K3Yg==",
"license": "Apache-2.0",
"dependencies": {
"@alicloud/tea-typescript": "^1.5.1",
"kitx": "^2.0.0"
}
},
"node_modules/@alicloud/gateway-pop": {
"version": "0.0.6",
"resolved": "https://registry.npmjs.org/@alicloud/gateway-pop/-/gateway-pop-0.0.6.tgz",
"integrity": "sha512-KF4I+JvfYuLKc3fWeWYIZ7lOVJ9jRW0sQXdXidZn1DKZ978ncfGf7i0LBfONGk4OxvNb/HD3/0yYhkgZgPbKtA==",
"license": "ISC",
"dependencies": {
"@alicloud/credentials": "^2",
"@alicloud/darabonba-array": "^0.1.0",
"@alicloud/darabonba-encode-util": "^0.0.2",
"@alicloud/darabonba-map": "^0.0.1",
"@alicloud/darabonba-signature-util": "^0.0.4",
"@alicloud/darabonba-string": "^1.0.2",
"@alicloud/endpoint-util": "^0.0.1",
"@alicloud/gateway-spi": "^0.0.8",
"@alicloud/openapi-util": "^0.3.2",
"@alicloud/tea-typescript": "^1.7.1",
"@alicloud/tea-util": "^1.4.8"
}
},
"node_modules/@alicloud/gateway-spi": {
"version": "0.0.8",
"resolved": "https://registry.npmjs.org/@alicloud/gateway-spi/-/gateway-spi-0.0.8.tgz",
"integrity": "sha512-KM7fu5asjxZPmrz9sJGHJeSU+cNQNOxW+SFmgmAIrITui5hXL2LB+KNRuzWmlwPjnuA2X3/keq9h6++S9jcV5g==",
"license": "ISC",
"dependencies": {
"@alicloud/credentials": "^2",
"@alicloud/tea-typescript": "^1.7.1"
}
},
"node_modules/@alicloud/openapi-core": {
"version": "1.0.8",
"resolved": "https://registry.npmjs.org/@alicloud/openapi-core/-/openapi-core-1.0.8.tgz",
"integrity": "sha512-xs8LdgMDcEUqv13kZ4nl+Vd+Fc1mixTF0g1lm5QSrNWDaLUUD5vqpuMtvUZnKNnq9PITfUQ+8KunAvNQJpFo8g==",
"hasInstallScript": true,
"license": "ISC",
"dependencies": {
"@alicloud/credentials": "^2.4.2",
"@alicloud/gateway-pop": "0.0.6",
"@alicloud/gateway-spi": "^0.0.8",
"@darabonba/typescript": "^1.0.5"
}
},
"node_modules/@alicloud/openapi-util": {
"version": "0.3.3",
"resolved": "https://registry.npmjs.org/@alicloud/openapi-util/-/openapi-util-0.3.3.tgz",
"integrity": "sha512-vf0cQ/q8R2U7ZO88X5hDiu1yV3t/WexRj+YycWxRutkH/xVXfkmpRgps8lmNEk7Ar+0xnY8+daN2T+2OyB9F4A==",
"license": "ISC",
"dependencies": {
"@alicloud/tea-typescript": "^1.7.1",
"@alicloud/tea-util": "^1.3.0",
"kitx": "^2.1.0",
"sm3": "^1.0.3"
}
},
"node_modules/@alicloud/tea-typescript": {
"version": "1.8.0",
"resolved": "https://registry.npmjs.org/@alicloud/tea-typescript/-/tea-typescript-1.8.0.tgz",
"integrity": "sha512-CWXWaquauJf0sW30mgJRVu9aaXyBth5uMBCUc+5vKTK1zlgf3hIqRUjJZbjlwHwQ5y9anwcu18r48nOZb7l2QQ==",
"license": "ISC",
"dependencies": {
"@types/node": "^12.0.2",
"httpx": "^2.2.6"
}
},
"node_modules/@alicloud/tea-typescript/node_modules/@types/node": {
"version": "12.20.55",
"resolved": "https://registry.npmjs.org/@types/node/-/node-12.20.55.tgz",
"integrity": "sha512-J8xLz7q2OFulZ2cyGTLE1TbbZcjpno7FaN6zdJNrgAdrJ+DZzh/uFR6YrTb4C+nXakvud8Q4+rbhoIWlYQbUFQ==",
"license": "MIT"
},
"node_modules/@alicloud/tea-util": {
"version": "1.4.11",
"resolved": "https://registry.npmjs.org/@alicloud/tea-util/-/tea-util-1.4.11.tgz",
"integrity": "sha512-HyPEEQ8F0WoZegiCp7sVdrdm6eBOB+GCvGl4182u69LDFktxfirGLcAx3WExUr1zFWkq2OSmBroTwKQ4w/+Yww==",
"license": "Apache-2.0",
"dependencies": {
"@alicloud/tea-typescript": "^1.5.1",
"@darabonba/typescript": "^1.0.0",
"kitx": "^2.0.0"
}
},
"node_modules/@anthropic-ai/claude-agent-sdk": {
"version": "0.3.202",
"resolved": "https://registry.npmjs.org/@anthropic-ai/claude-agent-sdk/-/claude-agent-sdk-0.3.202.tgz",
@@ -234,6 +406,24 @@
"node": ">=6.9.0"
}
},
"node_modules/@darabonba/typescript": {
"version": "1.0.5",
"resolved": "https://registry.npmjs.org/@darabonba/typescript/-/typescript-1.0.5.tgz",
"integrity": "sha512-pfxHFVM8I3h8K2o8skpDQLMmR5iAeQ2eNpP1HrKDEm/9ZPF8aKwDKPwwEszT1NnIo0Y3++E7x1YsVkVpGjA/xw==",
"license": "Apache License 2.0",
"dependencies": {
"@alicloud/tea-typescript": "^1.5.1",
"http-proxy-agent": "^5.0.0",
"https-proxy-agent": "^5.0.1",
"httpx": "^2.3.2",
"lodash": "^4.17.21",
"moment": "^2.30.1",
"moment-timezone": "^0.5.45",
"socks-proxy-agent": "^6.2.1",
"ws": "^8.18.0",
"xml2js": "^0.6.2"
}
},
"node_modules/@emnapi/core": {
"version": "1.11.2",
"resolved": "https://registry.npmjs.org/@emnapi/core/-/core-1.11.2.tgz",
@@ -1396,6 +1586,15 @@
"integrity": "sha512-l2aFy5jALhniG5HgqrD6jXLi/rUWrKvqN/qJx6yoJsgKhblVd+iqqU4RCXavm/jPityDo5TCvKMnpjKnOriy0w==",
"license": "MIT"
},
"node_modules/@tootallnate/once": {
"version": "2.0.1",
"resolved": "https://registry.npmjs.org/@tootallnate/once/-/once-2.0.1.tgz",
"integrity": "sha512-HqmEUIGRJ5fSXchkVgR5F7qn48bDBzv0kWj/Kfu5e6uci4UlEeng4331LnBkWffb++Ei3FOVLxo8JJWMFBDMeQ==",
"license": "MIT",
"engines": {
"node": ">= 10"
}
},
"node_modules/@tybys/wasm-util": {
"version": "0.10.3",
"resolved": "https://registry.npmjs.org/@tybys/wasm-util/-/wasm-util-0.10.3.tgz",
@@ -2830,6 +3029,20 @@
"url": "https://opencollective.com/express"
}
},
"node_modules/http-proxy-agent": {
"version": "5.0.0",
"resolved": "https://registry.npmjs.org/http-proxy-agent/-/http-proxy-agent-5.0.0.tgz",
"integrity": "sha512-n2hY8YdoRE1i7r6M0w9DIw5GgZN0G25P8zLCRQ8rjXtTU3vsNFBI/vWK/UIeE6g5MUUz6avwAPXmL6Fy9D/90w==",
"license": "MIT",
"dependencies": {
"@tootallnate/once": "2",
"agent-base": "6",
"debug": "4"
},
"engines": {
"node": ">= 6"
}
},
"node_modules/https-proxy-agent": {
"version": "5.0.1",
"resolved": "https://registry.npmjs.org/https-proxy-agent/-/https-proxy-agent-5.0.1.tgz",
@@ -2843,6 +3056,25 @@
"node": ">= 6"
}
},
"node_modules/httpx": {
"version": "2.3.3",
"resolved": "https://registry.npmjs.org/httpx/-/httpx-2.3.3.tgz",
"integrity": "sha512-k1qv94u1b6e+XKCxVbLgYlOypVP9MPGpnN5G/vxFf6tDO4V3xpz3d6FUOY/s8NtPgaq5RBVVgSB+7IHpVxMYzw==",
"license": "MIT",
"dependencies": {
"@types/node": "^20",
"debug": "^4.1.1"
}
},
"node_modules/httpx/node_modules/@types/node": {
"version": "20.19.43",
"resolved": "https://registry.npmjs.org/@types/node/-/node-20.19.43.tgz",
"integrity": "sha512-6oYBAi5ikg4Pl+kGsoYtawUMBT2zZMCvPNF7pVLnHZfd1zf38DRiWn/gT01RYCdUqkv7Fhr+C9ot4/tb+2sVvA==",
"license": "MIT",
"dependencies": {
"undici-types": "~6.21.0"
}
},
"node_modules/iconv-lite": {
"version": "0.7.3",
"resolved": "https://registry.npmjs.org/iconv-lite/-/iconv-lite-0.7.3.tgz",
@@ -2867,12 +3099,17 @@
"license": "ISC",
"peer": true
},
"node_modules/ini": {
"version": "1.3.8",
"resolved": "https://registry.npmjs.org/ini/-/ini-1.3.8.tgz",
"integrity": "sha512-JV/yugV2uzW5iMRSiZAyDtQd+nxtUnjeLt0acNdw98kKLrvuRVyB80tsREOE7yvGVgalhZ6RNXCmEHkUKBKxew==",
"license": "ISC"
},
"node_modules/ip-address": {
"version": "10.2.0",
"resolved": "https://registry.npmjs.org/ip-address/-/ip-address-10.2.0.tgz",
"integrity": "sha512-/+S6j4E9AHvW9SWMSEY9Xfy66O5PWvVEJ08O0y5JGyEKQpojb0K0GKpz/v5HJ/G0vi3D2sjGK78119oXZeE0qA==",
"license": "MIT",
"peer": true,
"engines": {
"node": ">= 12"
}
@@ -2972,6 +3209,15 @@
"license": "BSD-2-Clause",
"peer": true
},
"node_modules/kitx": {
"version": "2.2.0",
"resolved": "https://registry.npmjs.org/kitx/-/kitx-2.2.0.tgz",
"integrity": "sha512-tBMwe6AALTBQJb0woQDD40734NKzb0Kzi3k7wQj9ar3AbP9oqhoVrdXPh7rk2r00/glIgd0YbToIUJsnxWMiIg==",
"license": "MIT",
"dependencies": {
"@types/node": "^22.5.4"
}
},
"node_modules/light-my-request": {
"version": "6.6.0",
"resolved": "https://registry.npmjs.org/light-my-request/-/light-my-request-6.6.0.tgz",
@@ -3270,6 +3516,12 @@
"url": "https://opencollective.com/parcel"
}
},
"node_modules/lodash": {
"version": "4.18.1",
"resolved": "https://registry.npmjs.org/lodash/-/lodash-4.18.1.tgz",
"integrity": "sha512-dMInicTPVE8d1e5otfwmmjlxkZoUpiVLwyeTdUsi/Caj/gfzzblBcCE5sRHV/AsjuCmxWrte2TNGSYuCeCq+0Q==",
"license": "MIT"
},
"node_modules/lodash.identity": {
"version": "3.0.0",
"resolved": "https://registry.npmjs.org/lodash.identity/-/lodash.identity-3.0.0.tgz",
@@ -3357,6 +3609,27 @@
"node": ">= 0.6"
}
},
"node_modules/moment": {
"version": "2.30.1",
"resolved": "https://registry.npmjs.org/moment/-/moment-2.30.1.tgz",
"integrity": "sha512-uEmtNhbDOrWPFS+hdjFCBfy9f2YoyzRpwcl+DqpC6taX21FzsTLQVbMV/W7PzNSX6x/bhC1zA3c2UQ5NzH6how==",
"license": "MIT",
"engines": {
"node": "*"
}
},
"node_modules/moment-timezone": {
"version": "0.5.48",
"resolved": "https://registry.npmjs.org/moment-timezone/-/moment-timezone-0.5.48.tgz",
"integrity": "sha512-f22b8LV1gbTO2ms2j2z13MuPogNoh5UzxL3nzNAYKGraILnbGc9NEE6dyiiiLv46DGRb8A4kg8UKWLjPthxBHw==",
"license": "MIT",
"dependencies": {
"moment": "^2.29.4"
},
"engines": {
"node": "*"
}
},
"node_modules/ms": {
"version": "2.1.3",
"resolved": "https://registry.npmjs.org/ms/-/ms-2.1.3.tgz",
@@ -3976,6 +4249,15 @@
"license": "MIT",
"peer": true
},
"node_modules/sax": {
"version": "1.6.0",
"resolved": "https://registry.npmjs.org/sax/-/sax-1.6.0.tgz",
"integrity": "sha512-6R3J5M4AcbtLUdZmRv2SygeVaM7IhrLXu9BmnOGmmACak8fiUtOsYNWUS4uK7upbmHIBbLBeFeI//477BKLBzA==",
"license": "BlueOak-1.0.0",
"engines": {
"node": ">=11.0.0"
}
},
"node_modules/secure-json-parse": {
"version": "4.1.0",
"resolved": "https://registry.npmjs.org/secure-json-parse/-/secure-json-parse-4.1.0.tgz",
@@ -4193,6 +4475,50 @@
"dev": true,
"license": "ISC"
},
"node_modules/sm3": {
"version": "1.0.3",
"resolved": "https://registry.npmjs.org/sm3/-/sm3-1.0.3.tgz",
"integrity": "sha512-KyFkIfr8QBlFG3uc3NaljaXdYcsbRy1KrSfc4tsQV8jW68jAktGeOcifu530Vx/5LC+PULHT0Rv8LiI8Gw+c1g==",
"license": "MIT"
},
"node_modules/smart-buffer": {
"version": "4.2.0",
"resolved": "https://registry.npmjs.org/smart-buffer/-/smart-buffer-4.2.0.tgz",
"integrity": "sha512-94hK0Hh8rPqQl2xXc3HsaBoOXKV20MToPkcXvwbISWLEs+64sBq5kFgn2kJDHb1Pry9yrP0dxrCI9RRci7RXKg==",
"license": "MIT",
"engines": {
"node": ">= 6.0.0",
"npm": ">= 3.0.0"
}
},
"node_modules/socks": {
"version": "2.8.9",
"resolved": "https://registry.npmjs.org/socks/-/socks-2.8.9.tgz",
"integrity": "sha512-LJhUYUvItdQ0LkJTmPeaEObWXAqFyfmP85x0tch/ez9cahmhlBBLbIqDFnvBnUJGagb0JbIQrkBs1wJ+yRYpEw==",
"license": "MIT",
"dependencies": {
"ip-address": "^10.1.1",
"smart-buffer": "^4.2.0"
},
"engines": {
"node": ">= 10.0.0",
"npm": ">= 3.0.0"
}
},
"node_modules/socks-proxy-agent": {
"version": "6.2.1",
"resolved": "https://registry.npmjs.org/socks-proxy-agent/-/socks-proxy-agent-6.2.1.tgz",
"integrity": "sha512-a6KW9G+6B3nWZ1yB8G7pJwL3ggLy1uTzKAgCb7ttblwqdz9fMGJUuTy3uFzEP48FAs9FLILlmzDlE2JJhVQaXQ==",
"license": "MIT",
"dependencies": {
"agent-base": "^6.0.2",
"debug": "^4.3.3",
"socks": "^2.6.2"
},
"engines": {
"node": ">= 10"
}
},
"node_modules/sonic-boom": {
"version": "4.2.1",
"resolved": "https://registry.npmjs.org/sonic-boom/-/sonic-boom-4.2.1.tgz",
@@ -4700,6 +5026,28 @@
}
}
},
"node_modules/xml2js": {
"version": "0.6.2",
"resolved": "https://registry.npmjs.org/xml2js/-/xml2js-0.6.2.tgz",
"integrity": "sha512-T4rieHaC1EXcES0Kxxj4JWgaUQHDk+qwHcYOCFHfiwKz7tOVPLq7Hjq9dM1WCMhylqMEfP7hMcOIChvotiZegA==",
"license": "MIT",
"dependencies": {
"sax": ">=0.6.0",
"xmlbuilder": "~11.0.0"
},
"engines": {
"node": ">=4.0.0"
}
},
"node_modules/xmlbuilder": {
"version": "11.0.1",
"resolved": "https://registry.npmjs.org/xmlbuilder/-/xmlbuilder-11.0.1.tgz",
"integrity": "sha512-fDlsI/kFEx7gLvbecc0/ohLG50fugQp8ryHzMTuW9vSa1GJ0XYWKnhsUx7oie3G98+r56aTQIUB4kht42R3JvA==",
"license": "MIT",
"engines": {
"node": ">=4.0"
}
},
"node_modules/zod": {
"version": "4.4.3",
"resolved": "https://registry.npmjs.org/zod/-/zod-4.4.3.tgz",
+4 -1
View File
@@ -1,12 +1,15 @@
{
"name": "@paradigm/hub",
"version": "0.0.31",
"version": "0.0.32",
"private": true,
"type": "module",
"engines": {
"node": ">=24"
},
"dependencies": {
"@alicloud/credentials": "^2.4.5",
"@alicloud/docmind-api20220711": "^1.4.15",
"@alicloud/tea-util": "^1.4.11",
"@anthropic-ai/claude-agent-sdk": "^0.3.202",
"@fastify/cookie": "^11.0.2",
"@larksuiteoapi/node-sdk": "^1.70.0",
@@ -0,0 +1,64 @@
-- ADR-0026: append-only UsageFact ledger. AgentRun.costUsd/inputTokens/
-- outputTokens become a derived rollup cache; the truth is in UsageFact.
CREATE TABLE "UsageFact" (
"id" TEXT NOT NULL,
"runId" TEXT NOT NULL,
"occurredAt" TIMESTAMP(3) NOT NULL,
"kind" TEXT NOT NULL,
"provider" TEXT NOT NULL,
"model" TEXT,
"inputTokens" INTEGER,
"outputTokens" INTEGER,
"quantity" DECIMAL(18, 6),
"unit" TEXT,
"costUsd" DECIMAL(18, 8),
"costSource" TEXT NOT NULL,
"capabilityId" TEXT,
"correlationId" TEXT,
"metadata" JSONB NOT NULL,
"createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
CONSTRAINT "UsageFact_pkey" PRIMARY KEY ("id")
);
CREATE INDEX "UsageFact_runId_occurredAt_idx" ON "UsageFact"("runId", "occurredAt");
CREATE INDEX "UsageFact_runId_kind_idx" ON "UsageFact"("runId", "kind");
CREATE INDEX "UsageFact_provider_model_occurredAt_idx"
ON "UsageFact"("provider", "model", "occurredAt");
CREATE INDEX "UsageFact_capabilityId_occurredAt_idx"
ON "UsageFact"("capabilityId", "occurredAt");
ALTER TABLE "UsageFact"
ADD CONSTRAINT "UsageFact_runId_fkey"
FOREIGN KEY ("runId") REFERENCES "AgentRun"("id")
ON DELETE CASCADE ON UPDATE CASCADE;
-- Backfill: one synthetic model_completion fact per AgentRun that has any
-- recorded usage (cost or tokens). correlationId = run id marks these as
-- backfill artefacts (real facts use an external correlation id or null).
-- Runs with no recorded usage stay fact-less and remain runsWithoutCost,
-- matching ADR-0022's "missing cost ≠ zero" rule and the pre-existing
-- migration 20260709143000_agent_run_cost_tracking's "no backfill for the
-- truly unrecorded" stance.
INSERT INTO "UsageFact" (
"id", "runId", "occurredAt", "kind", "provider", "model",
"inputTokens", "outputTokens", "costUsd", "costSource",
"correlationId", "metadata"
)
SELECT
'usagefact_backfill_' || "AgentRun"."id",
"AgentRun"."id",
COALESCE("AgentRun"."finishedAt", "AgentRun"."startedAt", CURRENT_TIMESTAMP),
'model_completion',
"AgentRun"."provider",
"AgentRun"."model",
"AgentRun"."inputTokens",
"AgentRun"."outputTokens",
"AgentRun"."costUsd",
COALESCE("AgentRun"."costSource", 'unknown'),
"AgentRun"."id",
'{}'::jsonb
FROM "AgentRun"
WHERE "AgentRun"."costUsd" IS NOT NULL
OR "AgentRun"."inputTokens" IS NOT NULL
OR "AgentRun"."outputTokens" IS NOT NULL;
@@ -0,0 +1,69 @@
-- ADR-0027: org-scoped capability connections for external document/media
-- transforms (PDF→MD, audio/video→text, …). Mirrors the Feishu Application
-- Connection shape: reuses the ADR-0024 envelope machinery with its own
-- payload schema and readiness probe, distinct from the model-provider
-- connection.
CREATE TABLE "OrganizationCapabilityConnection" (
"id" TEXT NOT NULL,
"organizationId" TEXT NOT NULL,
"capabilityId" TEXT NOT NULL,
"status" "OrganizationConnectionStatus" NOT NULL DEFAULT 'DRAFT',
"activeSecretVersionId" TEXT,
"activatedAt" TIMESTAMP(3),
"disabledAt" TIMESTAMP(3),
"createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
"updatedAt" TIMESTAMP(3) NOT NULL,
CONSTRAINT "OrganizationCapabilityConnection_pkey" PRIMARY KEY ("id")
);
CREATE UNIQUE INDEX "OrganizationCapabilityConnection_organizationId_capabilityId_key"
ON "OrganizationCapabilityConnection"("organizationId", "capabilityId");
CREATE UNIQUE INDEX "OrganizationCapabilityConnection_activeSecretVersionId_key"
ON "OrganizationCapabilityConnection"("activeSecretVersionId");
CREATE INDEX "OrganizationCapabilityConnection_organizationId_status_idx"
ON "OrganizationCapabilityConnection"("organizationId", "status");
CREATE INDEX "OrganizationCapabilityConnection_capabilityId_status_idx"
ON "OrganizationCapabilityConnection"("capabilityId", "status");
CREATE TABLE "CapabilityCredentialVersion" (
"id" TEXT NOT NULL,
"connectionId" TEXT NOT NULL,
"version" INTEGER NOT NULL,
"envelopeVersion" INTEGER NOT NULL DEFAULT 1,
"keyId" TEXT NOT NULL,
"envelope" JSONB NOT NULL,
"createdByUserId" TEXT,
"createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
"retiredAt" TIMESTAMP(3),
CONSTRAINT "CapabilityCredentialVersion_pkey" PRIMARY KEY ("id")
);
CREATE UNIQUE INDEX "CapabilityCredentialVersion_connectionId_version_key"
ON "CapabilityCredentialVersion"("connectionId", "version");
CREATE INDEX "CapabilityCredentialVersion_connectionId_retiredAt_idx"
ON "CapabilityCredentialVersion"("connectionId", "retiredAt");
CREATE INDEX "CapabilityCredentialVersion_keyId_idx"
ON "CapabilityCredentialVersion"("keyId");
CREATE INDEX "CapabilityCredentialVersion_createdByUserId_idx"
ON "CapabilityCredentialVersion"("createdByUserId");
ALTER TABLE "OrganizationCapabilityConnection"
ADD CONSTRAINT "OrganizationCapabilityConnection_organizationId_fkey"
FOREIGN KEY ("organizationId") REFERENCES "Organization"("id")
ON DELETE CASCADE ON UPDATE CASCADE;
ALTER TABLE "OrganizationCapabilityConnection"
ADD CONSTRAINT "OrganizationCapabilityConnection_activeSecretVersionId_fkey"
FOREIGN KEY ("activeSecretVersionId") REFERENCES "CapabilityCredentialVersion"("id")
ON DELETE RESTRICT ON UPDATE CASCADE;
ALTER TABLE "CapabilityCredentialVersion"
ADD CONSTRAINT "CapabilityCredentialVersion_connectionId_fkey"
FOREIGN KEY ("connectionId") REFERENCES "OrganizationCapabilityConnection"("id")
ON DELETE CASCADE ON UPDATE CASCADE;
ALTER TABLE "CapabilityCredentialVersion"
ADD CONSTRAINT "CapabilityCredentialVersion_createdByUserId_fkey"
FOREIGN KEY ("createdByUserId") REFERENCES "User"("id")
ON DELETE SET NULL ON UPDATE CASCADE;
+108 -2
View File
@@ -44,6 +44,7 @@ model Organization {
externalDirectoryConnections ExternalDirectoryConnection[]
providerConnections OrganizationProviderConnection[]
feishuApplicationConnection OrganizationFeishuApplicationConnection?
capabilityConnections OrganizationCapabilityConnection[]
agentSkills OrganizationAgentSkill[]
agentRoles OrganizationAgentRole[]
projectGroupBindings ProjectGroupBinding[]
@@ -198,6 +199,7 @@ model User {
auditEntries AuditEntry[] @relation("auditActor")
providerCredentialVersions ProviderCredentialVersion[] @relation("providerCredentialVersionCreator")
feishuCredentialVersions FeishuApplicationCredentialVersion[] @relation("feishuCredentialVersionCreator")
capabilityCredentialVersions CapabilityCredentialVersion[] @relation("capabilityCredentialVersionCreator")
feishuIdentities FeishuUserIdentity[]
}
@@ -596,9 +598,13 @@ model AgentRun {
summary String?
inputTokens Int?
outputTokens Int?
/// Provider/gateway-reported cost in USD. Null means this run has no trusted cost fact.
/// ADR-0026: derived rollup cache of UsageFact rows for this run. The truth
/// is in UsageFact; this column is convenience for existing readers. Null
/// means this run has no trusted cost fact (ADR-0022: missing cost ≠ zero).
costUsd Decimal? @db.Decimal(18, 8)
/// Semantic source of costUsd, e.g. provider_reported. Not a pricing-estimate fallback.
/// ADR-0026: semantic source of the rolled-up costUsd (provider_reported |
/// pricebook_derived | unknown). Mirrors the dominant CostSource of the
/// run's facts; not a pricing-estimate fallback.
costSource String?
metadata Json
error String?
@@ -612,6 +618,7 @@ model AgentRun {
projectLock ProjectAgentLock?
messages AgentMessage[] @relation("runMessages")
fileChanges AgentFileChange[] @relation("runFileChanges")
usageFacts UsageFact[] @relation("runUsageFacts")
@@index([projectId, status])
@@index([projectId, finishedAt])
@@ -817,3 +824,102 @@ model AgentFileChange {
@@index([projectId])
@@index([path])
}
// --- Usage fact ledger (ADR-0026) ----------------------------------------
/// ADR-0026: one billable consumption event inside an AgentRun. Append-only;
/// the run's AgentRun.costUsd/inputTokens/outputTokens are a derived rollup
/// of these rows, not the truth. kind=external_capability carries capabilityId
/// for media transforms (PDF→MD, audio/video→text, …). costUsd=null means
/// unknown, not zero (ADR-0022). The kind set is OPEN: new kinds must be
/// surfaced, not silently folded in.
model UsageFact {
id String @id @default(cuid())
runId String
/// When the consumption happened. Pricebook derivation uses this, not
/// AgentRun.finishedAt, because an external capability may finish before
/// the run ends.
occurredAt DateTime
/// model_completion | external_capability | tool_proxy. OPEN set.
kind String
/// e.g. openrouter, mineru, openai_whisper.
provider String
model String?
inputTokens Int?
outputTokens Int?
/// Non-token meter (pages, audio_seconds, invocations). Coexists with tokens.
quantity Decimal? @db.Decimal(18, 6)
unit String?
/// USD. Null = unknown, NOT zero (ADR-0022). When null, costSource=unknown.
costUsd Decimal? @db.Decimal(18, 8)
/// provider_reported | pricebook_derived | unknown. Required even when
/// costUsd is null, so a reader can distinguish "reported zero" from
/// "no report at all".
costSource String
/// For kind=external_capability, the registered capability id
/// (e.g. pdf_to_md_bundle, audio_video_to_text). Null for model_completion.
capabilityId String?
/// External request id for reconciliation / idempotency. Not an aggregation key.
correlationId String?
metadata Json
createdAt DateTime @default(now())
run AgentRun @relation("runUsageFacts", fields: [runId], references: [id], onDelete: Cascade)
@@index([runId, occurredAt])
@@index([runId, kind])
@@index([provider, model, occurredAt])
@@index([capabilityId, occurredAt])
}
// --- External capability connections (ADR-0027) -------------------------
/// ADR-0027: org-scoped credential connection for an external capability
/// (PDF→MD, audio/video→text, …). Structurally mirrors the Feishu Application
/// Connection: reuses the ADR-0024 envelope (KEK→DEK→AES-256-GCM, AAD-bound)
/// but has its own payload schema and readiness probe, distinct from the
/// model-provider connection. Unique by (organizationId, capabilityId).
model OrganizationCapabilityConnection {
id String @id @default(cuid())
organizationId String
capabilityId String
status OrganizationConnectionStatus @default(DRAFT)
activeSecretVersionId String? @unique
activatedAt DateTime?
disabledAt DateTime?
createdAt DateTime @default(now())
updatedAt DateTime @updatedAt
organization Organization @relation(fields: [organizationId], references: [id], onDelete: Cascade)
secretVersions CapabilityCredentialVersion[] @relation("capabilityCredentialVersions")
activeSecretVersion CapabilityCredentialVersion? @relation("activeCapabilityCredentialVersion", fields: [activeSecretVersionId], references: [id], onDelete: Restrict)
@@unique([organizationId, capabilityId])
@@index([organizationId, status])
@@index([capabilityId, status])
}
/// ADR-0024/0027: one immutable authenticated envelope per capability secret
/// version. Same encryption machinery as Provider/Feishu credential versions;
/// the payload inside is CapabilitySecretPayloadV1 (baseUrl, apiToken,
/// optional projectId).
model CapabilityCredentialVersion {
id String @id @default(cuid())
connectionId String
version Int
envelopeVersion Int @default(1)
keyId String
envelope Json
createdByUserId String?
createdAt DateTime @default(now())
retiredAt DateTime?
connection OrganizationCapabilityConnection @relation("capabilityCredentialVersions", fields: [connectionId], references: [id], onDelete: Cascade)
activeFor OrganizationCapabilityConnection? @relation("activeCapabilityCredentialVersion")
createdBy User? @relation("capabilityCredentialVersionCreator", fields: [createdByUserId], references: [id], onDelete: SetNull)
@@unique([connectionId, version])
@@index([connectionId, retiredAt])
@@index([keyId])
@@index([createdByUserId])
}
+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;
}
@@ -0,0 +1,70 @@
/**
* ADR-0027: org-scoped capability credential resolver. Reuses the ADR-0024
* envelope decryption machinery (LocalSecretEnvelope) with
* purpose="capability", distinct from the model-provider and Feishu
* application connections. Fail-closed: no process-global fallback.
*
* Mirrors the Feishu application connection resolver shape, minus the
* readiness probe (capability probes are per-capability and injected by the
* adapter wiring, not this resolver).
*/
import type { PrismaClient } from "@prisma/client";
import { LocalSecretEnvelope, type SecretEnvelopeV1 } from "../security/secretEnvelope.js";
import {
CapabilityConnectionUnavailable,
type CapabilitySecretPayload,
type CapabilityId,
} from "./types.js";
const CAPABILITY_PURPOSE = "capability";
export interface ResolvedCapabilityCredential extends CapabilitySecretPayload {
readonly connectionId: string;
readonly organizationId: string;
readonly capabilityId: string;
}
/**
* Resolve the active capability credential for an organization. Throws
* CapabilityConnectionUnavailable when the org has no ACTIVE connection
* (fail-closed, ADR-0024). Decrypted plaintext exists only in the returned
* object for the duration of the capability call; it is never cached, logged,
* or passed to the Agent process.
*/
export async function resolveCapabilityCredential(
prisma: PrismaClient,
secrets: LocalSecretEnvelope,
input: { readonly organizationId: string; readonly capabilityId: CapabilityId },
): Promise<ResolvedCapabilityCredential> {
const connection = await prisma.organizationCapabilityConnection.findFirst({
where: {
organizationId: input.organizationId,
capabilityId: input.capabilityId,
status: "ACTIVE",
},
include: { activeSecretVersion: true },
});
if (connection === null || connection.activeSecretVersion === null) {
throw new CapabilityConnectionUnavailable(input.capabilityId, input.organizationId);
}
const version = connection.activeSecretVersion;
const binding = {
purpose: CAPABILITY_PURPOSE,
organizationId: connection.organizationId,
connectionId: connection.id,
secretVersionId: version.id,
};
const payload = secrets.decryptJson<CapabilitySecretPayload>(binding, version.envelope as unknown as SecretEnvelopeV1);
if (payload.schemaVersion !== 1) {
throw new Error(`unsupported capability secret schemaVersion: ${payload.schemaVersion}`);
}
return {
connectionId: connection.id,
organizationId: connection.organizationId,
capabilityId: connection.capabilityId,
schemaVersion: 1,
accessKeyId: payload.accessKeyId,
accessKeySecret: payload.accessKeySecret,
endpoint: payload.endpoint,
};
}
+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}`;
}
+231
View File
@@ -0,0 +1,231 @@
/**
* ADR-0027: Alibaba Cloud Document Mind (docmind) client.
*
* 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 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元/秒
*/
import $DocmindClient, {
SubmitDocParserJobAdvanceRequest,
QueryDocParserStatusRequest,
} from "@alicloud/docmind-api20220711";
import { RuntimeOptions } from "@alicloud/tea-util";
import { createReadStream } from "node:fs";
import { basename } from "node:path";
import type { CapabilitySecretPayload } from "./types.js";
/** A single extracted image downloaded from the markdown's OSS image URLs. */
export interface DocmindExtractedImage {
readonly filename: string;
readonly data: Uint8Array;
}
/** The structured result of parsing one document. */
export interface DocmindParseResult {
readonly markdown: string;
readonly images: readonly DocmindExtractedImage[];
readonly pageCount: number;
readonly costUsd: number | null;
readonly requestId: string | null;
}
export interface DocmindParseOptions {
readonly inputFilePath: string;
}
export interface CapabilityProviderClient {
parse(credential: CapabilitySecretPayload, options: DocmindParseOptions): Promise<DocmindParseResult>;
}
export class DocmindClientError extends Error {
constructor(
message: string,
readonly code: "docmind_unreachable" | "docmind_rejected" | "docmind_invalid_response" | "docmind_no_output" | "docmind_timeout",
readonly upstreamStatus?: number,
) {
super(message);
this.name = "DocmindClientError";
}
}
/** Cost per page in USD (0.04 CNY/page ≈ 0.0056 USD). */
const COST_PER_PAGE_USD = 0.0056;
const POLL_INTERVAL_MS = 10_000;
const POLL_TIMEOUT_MS = 5 * 60_000;
type DocmindConfig = ConstructorParameters<typeof $DocmindClient.default>[0];
export class AliyunDocmindClient implements CapabilityProviderClient {
async parse(credential: CapabilitySecretPayload, options: DocmindParseOptions): Promise<DocmindParseResult> {
const config: DocmindConfig = {
endpoint: credential.endpoint,
accessKeyId: credential.accessKeyId,
accessKeySecret: credential.accessKeySecret,
type: "access_key",
regionId: "cn-hangzhou",
} as DocmindConfig;
const client = new $DocmindClient.default(config);
// 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 fileStream = createReadStream(options.inputFilePath);
const advanceRequest = new SubmitDocParserJobAdvanceRequest({
fileUrlObject: fileStream,
fileName,
outputFormat: ["markdown"],
formulaEnhancement: true,
});
const runtime = new RuntimeOptions({});
let submitResponse;
try {
submitResponse = await client.submitDocParserJobAdvance(advanceRequest, runtime);
} catch (e) {
throw new DocmindClientError(
e instanceof Error ? e.message : String(e),
"docmind_unreachable",
);
}
const jobId = submitResponse.body?.data?.id;
if (jobId === undefined || jobId === null || jobId === "") {
throw new DocmindClientError("docmind submit returned no job id", "docmind_invalid_response");
}
// 2. Poll QueryDocParserStatus until data.status === "success" or "fail".
const deadline = Date.now() + POLL_TIMEOUT_MS;
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 {
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 (status !== "success") {
throw new DocmindClientError(`docmind job ${jobId} timed out (status: ${status})`, "docmind_timeout");
}
if (markdownUrl === null) {
throw new DocmindClientError("docmind returned no markdown output URL", "docmind_no_output");
}
// 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 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[] = [];
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 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.
}
}
// 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;
});
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 };
}
}
function sleep(ms: number): Promise<void> {
return new Promise((resolve) => {
setTimeout(resolve, ms);
});
}
/** 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`;
}
+152
View File
@@ -0,0 +1,152 @@
/**
* ADR-0027: pdf_to_md_bundle capability adapter.
*
* Converts a PDF in the run's workspace into a Markdown bundle (md + extracted
* images) by calling the MinerU document parsing service, writing outputs into
* the workspace, and recording consumption on a UsageFact (ADR-0026).
*
* Invariants (ADR-0027):
* 1. Credential isolation — the capability credential is resolved in Hub and
* never reaches the Agent process. The MineruClient receives it as a
* call argument, not from the environment.
* 2. Workspace containment — input and output paths are confined to the
* run's workspace dir (ADR-0018 AgentSurface). Escapes are rejected.
* 3. Mandatory fact — a successful invocation always writes ≥1 UsageFact
* with kind=external_capability, even when costUsd is null (ADR-0022:
* missing cost ≠ zero).
*/
import { mkdir, writeFile } from "node:fs/promises";
import { join, resolve, relative, isAbsolute } from "node:path";
import type { PrismaClient } from "@prisma/client";
import { LocalSecretEnvelope } from "../security/secretEnvelope.js";
import { resolveCapabilityCredential } from "./capabilityConnections.js";
import { DocmindClientError, type CapabilityProviderClient } from "./docmindClient.js";
import {
CAPABILITIES,
type CapabilityAdapter,
type CapabilityInvocationInput,
type CapabilityInvocationResult,
type CapabilityOutputArtifact,
} from "./types.js";
const CAPABILITY_ID = "pdf_to_md_bundle" as const;
const PROVIDER_ID = "aliyun_docmind";
/** Thrown when a requested path escapes the workspace root (ADR-0018). */
export class CapabilityPathEscape extends Error {
constructor(readonly requested: string, readonly workspaceDir: string) {
super(`capability path escapes workspace: ${requested} (root ${workspaceDir})`);
this.name = "CapabilityPathEscape";
}
}
/** Resolve a workspace-relative path, rejecting escapes (ADR-0018 AgentSurface). */
function confineToWorkspace(requestedPath: string, workspaceDir: string): string {
if (isAbsolute(requestedPath)) {
const rel = relative(workspaceDir, requestedPath);
if (rel.startsWith("..") || rel === "") {
throw new CapabilityPathEscape(requestedPath, workspaceDir);
}
return requestedPath;
}
const resolved = resolve(workspaceDir, requestedPath);
const rel = relative(workspaceDir, resolved);
if (rel.startsWith("..")) {
throw new CapabilityPathEscape(requestedPath, workspaceDir);
}
return resolved;
}
export interface PdfToMdBundleDeps {
readonly secrets: LocalSecretEnvelope;
readonly client: CapabilityProviderClient;
readonly prisma: PrismaClient;
}
/** Build the pdf_to_md_bundle adapter. The client is injectable for testing. */
export function createPdfToMdBundleAdapter(deps: PdfToMdBundleDeps): CapabilityAdapter {
return {
capabilityId: CAPABILITY_ID,
async invoke(input: CapabilityInvocationInput): Promise<CapabilityInvocationResult> {
const descriptor = CAPABILITIES[CAPABILITY_ID];
// 1. Resolve org-scoped credential (fail-closed, ADR-0024/0027).
const credential = await resolveCapabilityCredential(deps.prisma, deps.secrets, {
organizationId: input.organizationId,
capabilityId: CAPABILITY_ID,
});
// 2. Confine input + output paths to the workspace (ADR-0018).
const absoluteInput = confineToWorkspace(input.inputPath, input.workspaceDir);
const absoluteOutputDir = confineToWorkspace(input.outputDir, input.workspaceDir);
await mkdir(absoluteOutputDir, { recursive: true });
// 3. Call the backing service.
let result;
try {
result = await deps.client.parse(credential, { inputFilePath: absoluteInput });
} catch (e) {
if (e instanceof DocmindClientError) throw e;
throw new DocmindClientError(
e instanceof Error ? e.message : String(e),
"docmind_unreachable",
);
}
if (result.markdown === "") {
throw new DocmindClientError("docmind returned empty markdown", "docmind_no_output");
}
// 4. Write outputs into the workspace.
const artifacts: CapabilityOutputArtifact[] = [];
const mdPath = join(absoluteOutputDir, "document.md");
await writeFile(mdPath, result.markdown, "utf8");
artifacts.push({ path: relative(input.workspaceDir, mdPath), kind: "markdown" });
for (const image of result.images) {
const imagePath = join(absoluteOutputDir, image.filename);
const rel = relative(absoluteOutputDir, imagePath);
if (rel.startsWith("..")) {
// Defensive: image filename must not escape the output dir.
throw new CapabilityPathEscape(image.filename, absoluteOutputDir);
}
await writeFile(imagePath, image.data);
artifacts.push({ path: relative(input.workspaceDir, imagePath), kind: "image" });
}
// 5. Write the UsageFact (ADR-0026/0027). Always written on success;
// costUsd null means unknown, NOT zero (ADR-0022).
const occurredAt = new Date();
await deps.prisma.usageFact.create({
data: {
runId: input.runId,
occurredAt,
kind: "external_capability",
provider: PROVIDER_ID,
model: null,
inputTokens: null,
outputTokens: null,
quantity: result.pageCount,
unit: descriptor.meteringUnit,
costUsd: result.costUsd,
costSource: result.costUsd !== null ? "provider_reported" : "unknown",
capabilityId: CAPABILITY_ID,
correlationId: result.requestId,
metadata: {},
},
});
return {
artifacts,
consumption: {
provider: PROVIDER_ID,
model: null,
inputTokens: null,
outputTokens: null,
quantity: result.pageCount,
unit: descriptor.meteringUnit,
costUsd: result.costUsd,
correlationId: result.requestId,
},
};
},
};
}
+101
View File
@@ -0,0 +1,101 @@
/**
* ADR-0027: External capability types shared across the adapter layer.
*
* A capability is a platform-registered, org-enabled document/media transform
* invoked as a side effect of an AgentRun. The adapter resolves the org's
* active capability connection, calls the backing service via an injectable
* client, writes output into the run's workspace (AgentSurface, ADR-0018),
* and records consumption on a UsageFact (ADR-0026).
*/
import type { PrismaClient, Prisma } from "@prisma/client";
/** Stable capability identifiers registered with the platform (ADR-0027). */
export const CAPABILITY_IDS = [
"pdf_to_md_bundle",
"audio_video_to_text",
] as const;
export type CapabilityId = (typeof CAPABILITY_IDS)[number];
/** Non-token metering unit for a capability (ADR-0026/0027). */
export interface CapabilityDescriptor {
readonly id: CapabilityId;
readonly meteringUnit: string;
}
/** Known capabilities and their metering units. Code-level registry. */
export const CAPABILITIES: Readonly<Record<CapabilityId, CapabilityDescriptor>> = {
pdf_to_md_bundle: { id: "pdf_to_md_bundle", meteringUnit: "pages" },
audio_video_to_text: { id: "audio_video_to_text", meteringUnit: "audio_seconds" },
};
/** Input passed to a capability adapter invocation. */
export interface CapabilityInvocationInput {
readonly runId: string;
readonly organizationId: string;
readonly projectId: string;
/** Absolute workspace dir of the run's project (ADR-0018 surface root). */
readonly workspaceDir: string;
/** Workspace-relative path to the input file (PDF, audio, …). */
readonly inputPath: string;
/** Workspace-relative directory to write outputs into. Created if absent. */
readonly outputDir: string;
/** Prisma client for UsageFact writes. */
readonly prisma: PrismaClient;
}
/** A successfully produced output artifact (file written into workspace). */
export interface CapabilityOutputArtifact {
/** Workspace-relative path of the written artifact. */
readonly path: string;
readonly kind: "markdown" | "image" | "metadata" | "other";
}
/** Consumption recorded for one invocation (written to UsageFact). */
export interface CapabilityConsumption {
/** The backing service provider id (e.g. "mineru", "openai_whisper"). */
readonly provider: string;
/** Model id if the service reports one; null for non-model services. */
readonly model: string | null;
readonly inputTokens: number | null;
readonly outputTokens: number | null;
/** Non-token meter (page count, audio seconds). */
readonly quantity: number;
readonly unit: string;
/** USD cost if the service reported one; null = unknown (ADR-0022). */
readonly costUsd: number | null;
/** External request id for reconciliation. */
readonly correlationId: string | null;
}
/** Result of a successful capability invocation. */
export interface CapabilityInvocationResult {
readonly artifacts: readonly CapabilityOutputArtifact[];
readonly consumption: CapabilityConsumption;
}
/** A capability adapter: resolves credentials, calls the service, writes output. */
export interface CapabilityAdapter {
readonly capabilityId: CapabilityId;
invoke(input: CapabilityInvocationInput): Promise<CapabilityInvocationResult>;
}
/** Decrypted capability credential (CapabilitySecretPayloadV1, ADR-0027).
* Alibaba Cloud Document Mind (docmind) uses AccessKey ID + Secret + endpoint. */
export interface CapabilitySecretPayload {
readonly schemaVersion: 1;
readonly accessKeyId: string;
readonly accessKeySecret: string;
readonly endpoint: string;
}
/** Thrown when an org has no ACTIVE capability connection (fail-closed, ADR-0024). */
export class CapabilityConnectionUnavailable extends Error {
constructor(readonly capabilityId: string, readonly organizationId: string) {
super(`no ACTIVE capability connection for ${capabilityId} in org ${organizationId}`);
this.name = "CapabilityConnectionUnavailable";
}
}
/** Prisma transaction client type alias (for resolver signatures). */
export type TxClient = Prisma.TransactionClient;
+45 -15
View File
@@ -92,7 +92,18 @@ export function createSlashCommandRegistry(
finishedAt: { not: null },
},
orderBy: { finishedAt: "asc" },
select: { model: true, provider: true, inputTokens: true, outputTokens: true, costUsd: true },
select: {
usageFacts: {
select: {
provider: true,
model: true,
inputTokens: true,
outputTokens: true,
costUsd: true,
},
orderBy: { occurredAt: "asc" },
},
},
});
await sendText(rt, chatId, formatUsageReport(runs, scope), sendOptions);
},
@@ -118,18 +129,24 @@ async function currentRoleSessionIds(
return sessions.map((session) => session.id);
}
interface UsageRun {
readonly model: string;
interface UsageRunWithFacts {
readonly usageFacts: readonly {
readonly provider: string;
readonly model: string | null;
readonly inputTokens: number | null;
readonly outputTokens: number | null;
readonly costUsd: unknown;
}[];
}
function formatUsageReport(runs: readonly UsageRun[], scope: "current" | "project"): string {
function formatUsageReport(runs: readonly UsageRunWithFacts[], scope: "current" | "project"): string {
if (runs.length === 0) return `${scope === "current" ? "当前角色会话" : "当前项目"}还没有已结束的 Agent run。`;
// Bucket by (fact.provider, fact.model) — ADR-0026: external capabilities
// carry their own provider/model and contribute their own meter, so a run
// with a main loop + an external call lands in two buckets. A run with no
// cost-bearing fact is "unrecorded" (ADR-0022: missing cost ≠ zero).
const buckets = new Map<string, {
provider: string; model: string; runs: number; inputTokens: number; outputTokens: number; costUsd: number;
provider: string; model: string; facts: number; inputTokens: number; outputTokens: number; costUsd: number;
}>();
let recordedRuns = 0;
let unrecordedRuns = 0;
@@ -137,22 +154,35 @@ function formatUsageReport(runs: readonly UsageRun[], scope: "current" | "projec
let totalOutputTokens = 0;
let totalCostUsd = 0;
for (const run of runs) {
const costUsd = decimalToNumberOrNull(run.costUsd);
if (costUsd === null) { unrecordedRuns++; continue; }
const facts = run.usageFacts;
if (facts.length === 0 || !facts.some((f) => f.costUsd !== null && f.costUsd !== undefined)) {
unrecordedRuns++;
// Tokens from unrecorded runs still count toward totals (matches pre-0026).
for (const f of facts) {
totalInputTokens += f.inputTokens ?? 0;
totalOutputTokens += f.outputTokens ?? 0;
}
continue;
}
recordedRuns++;
const key = `${run.provider}\u0000${run.model}`;
for (const f of facts) {
const costUsd = decimalToNumberOrNull(f.costUsd);
if (costUsd === null) continue;
const model = f.model ?? "(unknown model)";
const key = `${f.provider}\u0000${model}`;
const bucket = buckets.get(key) ?? {
provider: run.provider, model: run.model, runs: 0, inputTokens: 0, outputTokens: 0, costUsd: 0,
provider: f.provider, model, facts: 0, inputTokens: 0, outputTokens: 0, costUsd: 0,
};
bucket.runs++;
bucket.inputTokens += run.inputTokens ?? 0;
bucket.outputTokens += run.outputTokens ?? 0;
bucket.facts += 1;
bucket.inputTokens += f.inputTokens ?? 0;
bucket.outputTokens += f.outputTokens ?? 0;
bucket.costUsd += costUsd;
buckets.set(key, bucket);
totalInputTokens += run.inputTokens ?? 0;
totalOutputTokens += run.outputTokens ?? 0;
totalInputTokens += f.inputTokens ?? 0;
totalOutputTokens += f.outputTokens ?? 0;
totalCostUsd += costUsd;
}
}
const lines = [
`${scope === "current" ? "当前角色会话" : "当前项目"}用量`,
`已记录成本: ${formatInteger(recordedRuns)} runs · ${formatUsd(totalCostUsd)}`,
@@ -160,7 +190,7 @@ function formatUsageReport(runs: readonly UsageRun[], scope: "current" | "projec
];
if (unrecordedRuns > 0) lines.push(`另有 ${formatInteger(unrecordedRuns)} 个 run 未记录成本。`);
for (const bucket of buckets.values()) {
lines.push(`- ${bucket.provider} / ${bucket.model}: ${formatInteger(bucket.runs)} runs, ${formatUsd(bucket.costUsd)}`);
lines.push(`- ${bucket.provider} / ${bucket.model}: ${formatInteger(bucket.facts)} , ${formatUsd(bucket.costUsd)}`);
}
return lines.join("\n");
}
+29 -2
View File
@@ -604,6 +604,32 @@ export function makeTriggerHandler(deps: TriggerDeps): TriggerHandler {
data: { metadata: mergeSessionMetadata(session.metadata, metadataPatch) },
});
}
// ADR-0026: the UsageFact ledger is the truth; AgentRun.costUsd /
// inputTokens / outputTokens are a derived rollup cache for existing
// readers. We write the fact first, then mirror it onto the run.
// They are separate statements (not one transaction) so the AgentRun
// row lock is held for the shortest possible window and concurrent
// workspace teardown cannot deadlock on the FK ShareLock. A crash
// between the two leaves the cache stale, but the cache is derived
// (the usage service re-reads facts), so staleness is recoverable;
// the reverse order would lose the truth and is unrecoverable.
const finishedAt = new Date();
const reportedCost = result.costUsd;
const costSource = reportedCost !== undefined ? "provider_reported" : "unknown";
await deps.prisma.usageFact.create({
data: {
runId: run.id,
occurredAt: finishedAt,
kind: "model_completion",
provider: providerId,
model,
inputTokens: result.usage.inputTokens,
outputTokens: result.usage.outputTokens,
costUsd: reportedCost ?? null,
costSource,
metadata: {},
},
});
await deps.prisma.agentRun.update({
where: { id: run.id },
data: {
@@ -618,9 +644,10 @@ export function makeTriggerHandler(deps: TriggerDeps): TriggerHandler {
: "FAILED",
inputTokens: result.usage.inputTokens,
outputTokens: result.usage.outputTokens,
...(result.costUsd !== undefined ? { costUsd: result.costUsd, costSource: "provider_reported" } : {}),
costUsd: reportedCost ?? null,
costSource,
error: result.error ?? null,
finishedAt: new Date(),
finishedAt,
},
});
await writeAudit(deps.prisma, {
+35
View File
@@ -73,9 +73,28 @@ export async function getSessionDetail(
inputTokens: true,
outputTokens: true,
costUsd: true,
costSource: true,
startedAt: true,
finishedAt: true,
error: true,
usageFacts: {
select: {
id: true,
occurredAt: true,
kind: true,
provider: true,
model: true,
inputTokens: true,
outputTokens: true,
quantity: true,
unit: true,
costUsd: true,
costSource: true,
capabilityId: true,
correlationId: true,
},
orderBy: { occurredAt: "asc" },
},
},
orderBy: { startedAt: "desc" },
take: 100,
@@ -106,9 +125,25 @@ export async function getSessionDetail(
inputTokens: r.inputTokens,
outputTokens: r.outputTokens,
costUsd: r.costUsd === null ? null : Number(r.costUsd),
costSource: r.costSource,
startedAt: r.startedAt.toISOString(),
finishedAt: r.finishedAt?.toISOString() ?? null,
error: r.error,
usageFacts: r.usageFacts.map((f) => ({
id: f.id,
occurredAt: f.occurredAt.toISOString(),
kind: f.kind,
provider: f.provider,
model: f.model,
inputTokens: f.inputTokens,
outputTokens: f.outputTokens,
quantity: f.quantity === null ? null : Number(f.quantity),
unit: f.unit,
costUsd: f.costUsd === null ? null : Number(f.costUsd),
costSource: f.costSource,
capabilityId: f.capabilityId,
correlationId: f.correlationId,
})),
})),
};
}
+87 -17
View File
@@ -1,8 +1,14 @@
/**
* Org / project usage rollups from AgentRun cost facts (ADR-0021).
* Org / project usage rollups from the UsageFact ledger (ADR-0021, ADR-0026).
* Operational usage accounting, not payment collection. Organizations may use
* either their own provider credentials or a distinct platform-managed
* connection; commercial billing remains outside the pilot scope.
*
* The truth is in `UsageFact`; `AgentRun.costUsd / inputTokens / outputTokens`
* are a derived rollup cache. This service reads `UsageFact` directly so that
* external-capability consumption (PDF→MD, ASR, …) is attributed correctly,
* not just the main model loop. A run with no cost-bearing fact is
* `runsWithoutCost` — ADR-0022: missing cost ≠ zero.
*/
import type { Prisma, PrismaClient } from "@prisma/client";
@@ -32,6 +38,53 @@ export interface UsageReport {
};
}
type FactRow = {
readonly inputTokens: number | null;
readonly outputTokens: number | null;
readonly costUsd: unknown;
};
type RunWithFacts = {
readonly projectId: string;
readonly usageFacts: readonly FactRow[];
};
/** A run is "recorded with cost" if any of its facts carries a known costUsd. */
function runHasRecordedCost(facts: readonly FactRow[]): boolean {
return facts.some((f) => f.costUsd !== null && f.costUsd !== undefined);
}
/** Decimal→number for Prisma Decimal; null/undefined → null. */
function factCostUsdToNumber(value: unknown): number | null {
if (value === null || value === undefined) return null;
const n = typeof value === "number" ? value : Number(value);
return Number.isFinite(n) ? n : null;
}
interface RunRollup {
readonly hasCost: boolean;
readonly inputTokens: number;
readonly outputTokens: number;
readonly costUsd: number | null;
}
function rollupRun(facts: readonly FactRow[]): RunRollup {
let inputTokens = 0;
let outputTokens = 0;
let costUsd: number | null = null;
let hasCost = false;
for (const f of facts) {
inputTokens += f.inputTokens ?? 0;
outputTokens += f.outputTokens ?? 0;
const c = factCostUsdToNumber(f.costUsd);
if (c !== null) {
hasCost = true;
costUsd = (costUsd ?? 0) + c;
}
}
return { hasCost, inputTokens, outputTokens, costUsd };
}
export async function getOrgUsage(
prisma: PrismaClient,
input: {
@@ -54,8 +107,9 @@ export async function getOrgUsage(
return emptyReport(input.from, input.to);
}
const projectIds = projects.map((p) => p.id);
const runWhere: Prisma.AgentRunWhereInput = {
projectId: { in: projects.map((p) => p.id) },
projectId: { in: projectIds },
...(input.from !== undefined || input.to !== undefined
? {
startedAt: {
@@ -66,20 +120,25 @@ export async function getOrgUsage(
: {}),
};
// Date range filters by run.startedAt, matching the pre-ADR-0026 semantics:
// a run is in or out of the window based on when it started, and all of its
// facts come along. Filtering facts by occurredAt independently would let a
// run contribute partial cost to a window it doesn't belong to.
const runs = await prisma.agentRun.findMany({
where: runWhere,
select: {
projectId: true,
inputTokens: true,
outputTokens: true,
costUsd: true,
usageFacts: {
select: { inputTokens: true, outputTokens: true, costUsd: true },
orderBy: { occurredAt: "asc" },
},
},
});
type MutableRow = {
projectId: string;
projectName: string;
folderId: string | null;
readonly projectId: string;
readonly projectName: string;
readonly folderId: string | null;
runCount: number;
runsWithCost: number;
runsWithoutCost: number;
@@ -103,22 +162,33 @@ export async function getOrgUsage(
});
}
for (const run of runs) {
for (const run of runs as readonly RunWithFacts[]) {
const row = byProject.get(run.projectId);
if (row === undefined) continue;
const roll = rollupRun(run.usageFacts);
row.runCount += 1;
row.inputTokens += run.inputTokens ?? 0;
row.outputTokens += run.outputTokens ?? 0;
if (run.costUsd === null) {
row.runsWithoutCost += 1;
} else {
if (roll.hasCost) {
row.runsWithCost += 1;
const cost = Number(run.costUsd);
row.costUsd = (row.costUsd ?? 0) + cost;
row.costUsd = (row.costUsd ?? 0) + (roll.costUsd ?? 0);
} else {
row.runsWithoutCost += 1;
}
row.inputTokens += roll.inputTokens;
row.outputTokens += roll.outputTokens;
}
const projectsOut: ProjectUsageRow[] = [...byProject.values()];
const projectsOut: ProjectUsageRow[] = [...byProject.values()].map((r) => ({
projectId: r.projectId,
projectName: r.projectName,
folderId: r.folderId,
runCount: r.runCount,
runsWithCost: r.runsWithCost,
runsWithoutCost: r.runsWithoutCost,
inputTokens: r.inputTokens,
outputTokens: r.outputTokens,
costUsd: r.costUsd,
}));
let runCount = 0;
let runsWithCost = 0;
let runsWithoutCost = 0;
@@ -0,0 +1,285 @@
import { mkdtemp, mkdir, writeFile, readFile, rm } from "node:fs/promises";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { prisma, resetDb, seedTestOrganization, testSecretEnvelope, DEFAULT_ORG_ID } from "./helpers.js";
import { createPdfToMdBundleAdapter, CapabilityPathEscape } from "../../src/capability/pdfToMdBundle.js";
import { DocmindClientError, type CapabilityProviderClient, type DocmindParseResult } from "../../src/capability/docmindClient.js";
import type { CapabilitySecretPayload } from "../../src/capability/types.js";
const CAPABILITY_ID = "pdf_to_md_bundle";
const PROVIDER_ID = "aliyun_docmind";
/**
* ADR-0027: pdf_to_md_bundle adapter contract tests. Proves:
* - workspace containment (input + output confined to run workspace)
* - fail-closed credential resolution (no ACTIVE connection → error)
* - UsageFact attribution with non-token meter (pages), costSource rules
* - outputs (md + images) written into workspace
* - docmind client errors propagate
* The client is mocked; the prisma + envelope + filesystem are real.
*/
describe("pdf_to_md_bundle capability adapter (ADR-0027)", () => {
let workspaceRoot: string;
let runId: string;
beforeEach(async () => {
await resetDb();
await seedTestOrganization();
workspaceRoot = await mkdtemp(join(tmpdir(), "cph-cap-"));
runId = "run-cap-test";
// AgentRun is required for the UsageFact FK.
await prisma.project.create({
data: {
id: "proj-cap",
organizationId: DEFAULT_ORG_ID,
name: "Cap Test",
workspaceDir: workspaceRoot,
},
});
await prisma.agentRun.create({
data: {
id: runId,
projectId: "proj-cap",
entrypoint: "FEISHU",
provider: "openrouter",
model: "mock-model",
status: "ACTIVE",
prompt: "convert this pdf",
metadata: {},
},
});
});
afterEach(async () => {
await rm(workspaceRoot, { recursive: true, force: true }).catch(() => {});
});
async function seedActiveCapabilityConnection(): Promise<void> {
const payload: CapabilitySecretPayload = {
schemaVersion: 1,
accessKeyId: "LTAI-test-key-id",
accessKeySecret: "test-secret-never-log",
endpoint: "docmind-api.cn-hangzhou.aliyuncs.com",
};
const connection = await prisma.organizationCapabilityConnection.create({
data: {
id: "cap-conn-1",
organizationId: DEFAULT_ORG_ID,
capabilityId: CAPABILITY_ID,
status: "ACTIVE",
activatedAt: new Date(),
},
});
const envelope = testSecretEnvelope.encryptJson(
{
purpose: "capability",
organizationId: DEFAULT_ORG_ID,
connectionId: connection.id,
secretVersionId: "cap-sv-1",
},
payload,
);
await prisma.capabilityCredentialVersion.create({
data: {
id: "cap-sv-1",
connectionId: connection.id,
version: 1,
envelope: envelope as object,
keyId: envelope.keyId,
},
});
await prisma.organizationCapabilityConnection.update({
where: { id: connection.id },
data: { activeSecretVersionId: "cap-sv-1" },
});
}
function mockDocmindClient(result: Partial<DocmindParseResult> = {}): CapabilityProviderClient {
const full: DocmindParseResult = {
markdown: "# Parsed Document\n\nHello world.\n\n![fig](page_1.jpg)\n",
images: [
{ filename: "page_1.jpg", data: new Uint8Array([0xff, 0xd8, 0xff, 0xe0]) },
],
pageCount: 3,
costUsd: 0.0168,
requestId: "docmind-job-abc",
...result,
};
return {
parse: vi.fn(async (): Promise<DocmindParseResult> => full),
};
}
function makeAdapter(client: CapabilityProviderClient) {
return createPdfToMdBundleAdapter({
secrets: testSecretEnvelope,
client,
prisma,
});
}
async function seedInputPdf(name: string = "input.pdf"): Promise<string> {
const dir = join(workspaceRoot, "sources");
await mkdir(dir, { recursive: true });
const path = join(dir, name);
await writeFile(path, "%PDF-1.4 fake pdf bytes");
return join("sources", name);
}
it("writes md + images into the workspace and records a UsageFact", async () => {
await seedActiveCapabilityConnection();
const inputPath = await seedInputPdf();
const client = mockDocmindClient();
const adapter = makeAdapter(client);
const result = await adapter.invoke({
runId,
organizationId: DEFAULT_ORG_ID,
projectId: "proj-cap",
workspaceDir: workspaceRoot,
inputPath,
outputDir: "output",
prisma,
});
// Outputs landed in workspace.
const md = await readFile(join(workspaceRoot, "output", "document.md"), "utf8");
expect(md).toContain("# Parsed Document");
expect(result.artifacts.some((a) => a.kind === "markdown" && a.path === "output/document.md")).toBe(true);
expect(result.artifacts.some((a) => a.kind === "image" && a.path === "output/page_1.jpg")).toBe(true);
// Client received the decrypted credential, not the env.
const call = (client.parse as ReturnType<typeof vi.fn>).mock.calls[0]?.[0] as CapabilitySecretPayload;
expect(call.accessKeyId).toBe("LTAI-test-key-id");
expect(call.endpoint).toBe("docmind-api.cn-hangzhou.aliyuncs.com");
// UsageFact written with non-token meter and provider_reported cost.
const facts = await prisma.usageFact.findMany({ where: { runId } });
expect(facts).toHaveLength(1);
const fact = facts[0]!;
expect(fact.kind).toBe("external_capability");
expect(fact.provider).toBe(PROVIDER_ID);
expect(fact.capabilityId).toBe(CAPABILITY_ID);
expect(Number(fact.quantity)).toBe(3);
expect(fact.unit).toBe("pages");
expect(Number(fact.costUsd!)).toBeCloseTo(0.0168, 8);
expect(fact.costSource).toBe("provider_reported");
expect(fact.correlationId).toBe("docmind-job-abc");
});
it("writes UsageFact with costSource=unknown when docmind reports no cost", async () => {
await seedActiveCapabilityConnection();
const inputPath = await seedInputPdf();
const client = mockDocmindClient({ costUsd: null, requestId: null });
const adapter = makeAdapter(client);
await adapter.invoke({
runId,
organizationId: DEFAULT_ORG_ID,
projectId: "proj-cap",
workspaceDir: workspaceRoot,
inputPath,
outputDir: "output",
prisma,
});
const fact = (await prisma.usageFact.findFirstOrThrow({ where: { runId } }));
expect(fact.costUsd).toBeNull();
expect(fact.costSource).toBe("unknown");
// quantity still recorded — missing cost ≠ zero consumption (ADR-0022).
expect(Number(fact.quantity)).toBe(3);
});
it("fails closed when the org has no ACTIVE capability connection", async () => {
// No connection seeded.
const inputPath = await seedInputPdf();
const adapter = makeAdapter(mockDocmindClient());
await expect(adapter.invoke({
runId,
organizationId: DEFAULT_ORG_ID,
projectId: "proj-cap",
workspaceDir: workspaceRoot,
inputPath,
outputDir: "output",
prisma,
})).rejects.toThrow(/no ACTIVE capability connection/);
// No UsageFact written for a failed resolution.
const facts = await prisma.usageFact.findMany({ where: { runId } });
expect(facts).toHaveLength(0);
});
it("rejects an input path that escapes the workspace", async () => {
await seedActiveCapabilityConnection();
const adapter = makeAdapter(mockDocmindClient());
await expect(adapter.invoke({
runId,
organizationId: DEFAULT_ORG_ID,
projectId: "proj-cap",
workspaceDir: workspaceRoot,
inputPath: "../../../etc/passwd",
outputDir: "output",
prisma,
})).rejects.toBeInstanceOf(CapabilityPathEscape);
});
it("rejects an output path that escapes the workspace", async () => {
await seedActiveCapabilityConnection();
const inputPath = await seedInputPdf();
const adapter = makeAdapter(mockDocmindClient());
await expect(adapter.invoke({
runId,
organizationId: DEFAULT_ORG_ID,
projectId: "proj-cap",
workspaceDir: workspaceRoot,
inputPath,
outputDir: "../../outside",
prisma,
})).rejects.toBeInstanceOf(CapabilityPathEscape);
});
it("propagates docmind client errors without writing a UsageFact", async () => {
await seedActiveCapabilityConnection();
const inputPath = await seedInputPdf();
const client: CapabilityProviderClient = {
parse: vi.fn(async (): Promise<DocmindParseResult> => {
throw new DocmindClientError("upstream 502", "docmind_rejected", 502);
}),
};
const adapter = makeAdapter(client);
await expect(adapter.invoke({
runId,
organizationId: DEFAULT_ORG_ID,
projectId: "proj-cap",
workspaceDir: workspaceRoot,
inputPath,
outputDir: "output",
prisma,
})).rejects.toThrow(/upstream 502/);
const facts = await prisma.usageFact.findMany({ where: { runId } });
expect(facts).toHaveLength(0);
});
it("rejects empty markdown output as a parse failure", async () => {
await seedActiveCapabilityConnection();
const inputPath = await seedInputPdf();
const client = mockDocmindClient({ markdown: "" });
const adapter = makeAdapter(client);
await expect(adapter.invoke({
runId,
organizationId: DEFAULT_ORG_ID,
projectId: "proj-cap",
workspaceDir: workspaceRoot,
inputPath,
outputDir: "output",
prisma,
})).rejects.toThrow(/empty markdown/);
});
});
+191
View File
@@ -0,0 +1,191 @@
import { describe, it, expect, beforeEach } from "vitest";
import { prisma, resetDb, seedTestOrganization, DEFAULT_ORG_ID } from "./helpers.js";
import { getOrgUsage, getProjectUsage } from "../../src/org/usage.js";
/**
* ADR-0026: org/project usage rolls up from the UsageFact ledger, not from the
* AgentRun scalar cache. These tests pin the fact-based aggregation contract:
* - a run with a cost-bearing fact → runsWithCost, costUsd summed
* - a run whose facts all have null costUsd → runsWithoutCost (≠ zero)
* - non-token meter (quantity/unit) is stored but does not corrupt token sums
* - multiple facts per run (main loop + external capability) are summed together
* - the AgentRun rollup cache is NOT consulted by getOrgUsage
*/
describe("usage rollup from UsageFact (ADR-0026)", () => {
beforeEach(async () => {
await resetDb();
await seedTestOrganization();
});
async function seedRun(opts: {
readonly projectId: string;
readonly projectName: string;
readonly folderId?: string;
readonly facts: ReadonlyArray<{
readonly kind?: string;
readonly provider?: string;
readonly model?: string;
readonly inputTokens?: number;
readonly outputTokens?: number;
readonly costUsd?: number | null;
readonly costSource?: string;
readonly capabilityId?: string;
readonly quantity?: number | null;
readonly unit?: string | null;
}>;
}): Promise<void> {
const project = await prisma.project.create({
data: {
id: opts.projectId,
organizationId: DEFAULT_ORG_ID,
name: opts.projectName,
workspaceDir: `/tmp/test-${opts.projectId}`,
...(opts.folderId !== undefined ? { folderId: opts.folderId } : {}),
},
});
const startedAt = new Date("2026-07-18T10:00:00Z");
const finishedAt = new Date("2026-07-18T10:05:00Z");
const run = await prisma.agentRun.create({
data: {
projectId: project.id,
entrypoint: "FEISHU",
provider: "openrouter",
model: "mock-model",
metadata: {},
status: "COMPLETED",
prompt: "test",
startedAt,
finishedAt,
// Deliberately set the rollup cache to values that differ from the
// facts, to prove getOrgUsage reads facts, not the cache.
inputTokens: 999,
outputTokens: 999,
costUsd: 999,
costSource: "provider_reported",
},
});
for (const [i, f] of opts.facts.entries()) {
await prisma.usageFact.create({
data: {
runId: run.id,
occurredAt: new Date(startedAt.getTime() + i * 1000),
kind: f.kind ?? "model_completion",
provider: f.provider ?? "openrouter",
model: f.model ?? "mock-model",
inputTokens: f.inputTokens ?? null,
outputTokens: f.outputTokens ?? null,
costUsd: f.costUsd ?? null,
costSource: f.costSource ?? (f.costUsd !== undefined && f.costUsd !== null ? "provider_reported" : "unknown"),
capabilityId: f.capabilityId ?? null,
quantity: f.quantity ?? null,
unit: f.unit ?? null,
metadata: {},
},
});
}
}
it("sums cost and tokens from facts, ignoring the AgentRun rollup cache", async () => {
await seedRun({
projectId: "proj-a",
projectName: "A",
facts: [{ inputTokens: 10, outputTokens: 5, costUsd: 0.0023 }],
});
await seedRun({
projectId: "proj-b",
projectName: "B",
facts: [{ inputTokens: 20, outputTokens: 10, costUsd: 0.01 }],
});
const report = await getOrgUsage(prisma, { organizationId: DEFAULT_ORG_ID });
expect(report.totals.runCount).toBe(2);
expect(report.totals.runsWithCost).toBe(2);
expect(report.totals.runsWithoutCost).toBe(0);
expect(report.totals.inputTokens).toBe(30);
expect(report.totals.outputTokens).toBe(15);
expect(report.totals.costUsd).toBeCloseTo(0.0123, 8);
});
it("treats a run with only null-cost facts as runsWithoutCost, never zero", async () => {
await seedRun({
projectId: "proj-unknown",
projectName: "Unknown",
facts: [{ inputTokens: 7, outputTokens: 3, costUsd: null, costSource: "unknown" }],
});
await seedRun({
projectId: "proj-known",
projectName: "Known",
facts: [{ inputTokens: 4, outputTokens: 2, costUsd: 0.05 }],
});
const report = await getOrgUsage(prisma, { organizationId: DEFAULT_ORG_ID });
expect(report.totals.runCount).toBe(2);
expect(report.totals.runsWithCost).toBe(1);
expect(report.totals.runsWithoutCost).toBe(1);
// The unknown-cost run still contributes its tokens.
expect(report.totals.inputTokens).toBe(11);
expect(report.totals.outputTokens).toBe(5);
// And its cost is NOT counted as zero — only the known cost is summed.
expect(report.totals.costUsd).toBeCloseTo(0.05, 8);
});
it("sums multiple facts per run (main loop + external capability)", async () => {
await seedRun({
projectId: "proj-multi",
projectName: "Multi",
facts: [
{ kind: "model_completion", inputTokens: 100, outputTokens: 50, costUsd: 0.02 },
{
kind: "external_capability",
provider: "mineru",
model: null,
inputTokens: null,
outputTokens: null,
costUsd: 0.015,
costSource: "provider_reported",
capabilityId: "pdf_to_md_bundle",
quantity: 12,
unit: "pages",
},
],
});
const report = await getOrgUsage(prisma, { organizationId: DEFAULT_ORG_ID });
expect(report.totals.runCount).toBe(1);
expect(report.totals.runsWithCost).toBe(1);
expect(report.totals.inputTokens).toBe(100);
expect(report.totals.outputTokens).toBe(50);
expect(report.totals.costUsd).toBeCloseTo(0.035, 8);
});
it("a run with no facts at all is runsWithoutCost and contributes nothing", async () => {
await seedRun({ projectId: "proj-empty", projectName: "Empty", facts: [] });
const report = await getOrgUsage(prisma, { organizationId: DEFAULT_ORG_ID });
expect(report.totals.runCount).toBe(1);
expect(report.totals.runsWithCost).toBe(0);
expect(report.totals.runsWithoutCost).toBe(1);
expect(report.totals.inputTokens).toBe(0);
expect(report.totals.outputTokens).toBe(0);
expect(report.totals.costUsd).toBeNull();
});
it("getProjectUsage returns just that project's row", async () => {
await seedRun({ projectId: "proj-x", projectName: "X", facts: [{ costUsd: 0.1, inputTokens: 1, outputTokens: 1 }] });
await seedRun({ projectId: "proj-y", projectName: "Y", facts: [{ costUsd: 0.2, inputTokens: 2, outputTokens: 2 }] });
const row = await getProjectUsage(prisma, { organizationId: DEFAULT_ORG_ID, projectId: "proj-x" });
expect(row.projectId).toBe("proj-x");
expect(row.runCount).toBe(1);
expect(row.costUsd).toBeCloseTo(0.1, 8);
});
it("returns an empty report for an org with no projects", async () => {
const report = await getOrgUsage(prisma, { organizationId: DEFAULT_ORG_ID });
expect(report.projects).toHaveLength(0);
expect(report.totals.runCount).toBe(0);
expect(report.totals.costUsd).toBeNull();
});
});
+3
View File
@@ -290,6 +290,9 @@ function mockPrisma(): PrismaClient {
findUnique: vi.fn(async () => null),
update: vi.fn(async () => ({ id: "run-1" })),
},
usageFact: {
create: vi.fn(async () => ({ id: "fact-1" })),
},
projectAgentLock: {
findUnique: vi.fn(async () => (lock === null ? null : { runId: lock.runId })),
create: vi.fn(async (args: { data: { projectId: string; runId: string } }) => {
+9 -1
View File
@@ -10,6 +10,8 @@ import Spec.System.Agent.Run
import Spec.System.Agent.AgentRole
import Spec.System.Agent.Memory
import Spec.System.Agent.AgentSurface
import Spec.System.Agent.Usage
import Spec.System.Agent.Capability
import Spec.System.Lock
import Spec.System.Permission
import Spec.System.PermissionGrant
@@ -39,11 +41,17 @@ likec4 已画出结构;这里补语义:
- `Memory` —— 按需上下文:锚点类别(ADR-0003)+ MCP 工具按 run/project 上下文授权。
- `AgentSurface` —— agent 执行面被 run 的工作区所界定(ADR-0018);与 Lock 正交——
Lock 限定并发,Surface 限定波及面。机制 OPEN。
- `Usage` —— Run 内 append-only 用量计量事实账本(ADR-0026);主模型 completion、
外部能力调用、代理网关旁路共用同一账本。fact 不持锁、不跨 run;`costUsd = none`
表示未知而非零(ADR-0022)。Run 上的 cost/token 标量是其 rollup cache。
- `Capability` —— 外部能力(PDF→MD、音视频→文本等)注册与调用(ADR-0027);
org-scoped 凭证连接复用 ADR-0024 信封但独立于 model-provider;调用是 Run 内副作用,
产物落 workspace,消耗记 UsageFact。capability credential 不进 Agent 进程。
- `Permission` —— read⊂edit⊂manage 角色体系、能力推导、单调性;force-release 在格外。
- `PermissionGrant` —— grant(resource×principal×role)与 settings(六 policy 旋钮)
(ADR-0004);组合规则 OPEN。
- `Audit` —— customer Project/Run 审计从简(内容 OPEN);Platform Audit 由
`PlatformAdministration` 独立承载。
标识符见 `Spec.Prelude`。决策出处:ADR-0001..0004, 0018, 0020..0024
标识符见 `Spec.Prelude`。决策出处:ADR-0001..0004, 0018, 0020..0027
-/
+104
View File
@@ -0,0 +1,104 @@
import Spec.Prelude
import Spec.System.Agent.Run
import Spec.System.Agent.Usage
/-!
# Capability —— 外部能力注册与调用(ADR-0027)
`AgentRun` 内可能调用**外部能力**:PDF→MD bundle、音视频→文本、OCR 等。这些能力
不是主 agent loop 的模型 completion,不持项目锁,不占 session 语义(ADR-0026 已拒绝
nested Run)。它们是 Run 内的副作用:读 workspace 输入、调外部服务、产物落 workspace、
写一条 `UsageFact`。
本模块 pin 三件事:
1. **ExternalCapability** —— 平台注册的、org 启用的文档/媒体转换服务,由稳定
`capabilityId` 标识。
2. **CapabilityConnection** —— org-scoped 凭证连接,复用 ADR-0024 信封机制但独立于
model-provider connection。只有 `active` connection 可被 resolver 使用。
3. **CapabilityInvocation 良构** —— 调用的输入/输出必须落在 run 的 workspace 内
(ADR-0018 `AgentSurface`),且调用产生的消耗必须以 `UsageFact` 记录。
凭证不经 Agent 子进程:Agent 只通过 capability adapter 间接调用,adapter 在 Hub 侧
解析 org-scoped 凭证并调用外部服务(ADR-0024 resolver boundary)。这与 model-provider
的 loopback proxy 是不同的机制——capability 调用不走 agent 的网络面,走 Hub 直连。
不变式(`PINNED`, ADR-0027):
- **凭证隔离** —— capability credential 不进 Agent 进程环境;adapter 在 Hub 侧解析。
- **workspace 边界** —— 输入路径与输出目录都必须落在 run 的 workspace 内(ADR-0018)。
- **必记 fact** —— 一次成功调用必须产生 ≥1 条 `UsageFact`(`kind = externalCapability`),
即使 `costUsd = none`(unknown ≠ 零,ADR-0022/0026)。
`CapabilityInvocation` 作为一等持久记录(status/retries/partial output)在 pilot 不做;
`UsageFact` 的 `capabilityId + correlationId` 是唯一可追溯痕迹(ADR-0027 Deferred)。
-/
namespace Spec.System.Agent
variable (I : Identifiers) (Path : Type)
/-- 外部能力的输入种类(`PINNED`, ADR-0027)。集合 `OPEN`——新增 kind 须 surface。 -/
inductive CapabilityInputKind where
/-- PDF 文档(`PINNED`)。 -/
| pdf
/-- 图片(`PINNED`)。 -/
| image
/-- 音频(`PINNED`)。 -/
| audio
/-- 视频(`PINNED`)。 -/
| video
/-- 外部能力(`PINNED`, ADR-0027)。平台注册的文档/媒体转换服务,由 org 启用。
一个 capability 由稳定 `capabilityId` 标识(如 `pdf_to_md_bundle`),声明接受的输入
种类与计量单位。org 通过 `CapabilityConnection` 启用它;adapter 是平台侧实现。 -/
structure ExternalCapability where
/-- 稳定能力标识(`PINNED`;如 `pdf_to_md_bundle`、`audio_video_to_text`)。 -/
capabilityId : String
/-- 接受的输入种类(`PINNED`, ADR-0027)。 -/
inputKind : CapabilityInputKind
/-- 计量单位(`OPEN`;如 `"pages"`、`"audio_seconds"`)。与 UsageFact.unit 对齐。 -/
meteringUnit : String
/-- capability connection 的运行态(`PINNED`, ADR-0024/0027):与 model-provider /
Feishu connection 同构,只有 `active` 可被 resolver 使用。 -/
inductive CapabilityConnectionStatus where
| draft
| active
| disabled
/-- Organization 的外部能力凭证连接(`PINNED`, ADR-0027)。复用 ADR-0024 信封机制
但独立于 model-provider connection:capability 服务(如 MinerU、Whisper)有自己的 auth
形状与 readiness probe,不走 OpenRouter `/v1/models`。唯一性键为
`(organization, capabilityId)`。 -/
structure OrganizationCapabilityConnection where
/-- connection 所属 organization(`PINNED`, ADR-0027)。 -/
organization : I.OrganizationId
/-- 该 connection 启用的 capability(`PINNED`, ADR-0027)。 -/
capability : ExternalCapability
/-- 运行态(`PINNED`, ADR-0024/0027)。 -/
status : CapabilityConnectionStatus
/-- 一次 capability 调用的良构约束(`PINNED`, ADR-0027)。输入路径与输出目录都必须落在
发起 run 的 workspace 内(ADR-0018 `AgentSurface`);`runWorkspace` 与 `pathWithin`
由平台提供(表示 `OPEN`)。逃逸即越权,拒绝。 -/
structure CapabilityInvocation where
/-- 发起调用的 run(`PINNED`;调用是 Run 内副作用,不跨 run,ADR-0026/0027)。 -/
run : I.RunId
/-- 被调用的 capability(`PINNED`)。 -/
capability : ExternalCapability
/-- 输入路径(workspace 内,`PINNED`, ADR-0018)。 -/
inputPath : Path
/-- 输出目录(workspace 内,`PINNED`, ADR-0018)。 -/
outputDir : Path
/-- capability 调用良构:输入与输出路径都落在 run 的 workspace 内(`PINNED`,
ADR-0018/0027)。这是 `AgentFileOp.Authorized` 在 capability 调用上的对应物。 -/
def CapabilityInvocation.Authorized
(inv : CapabilityInvocation I Path)
(runWorkspace : I.RunId Option Path)
(pathWithin : Path Path Prop) : Prop :=
w, runWorkspace inv.run = some w pathWithin inv.inputPath w pathWithin inv.outputDir w
end Spec.System.Agent
+92
View File
@@ -0,0 +1,92 @@
import Spec.Prelude
/-!
# Usage —— Run 内用量计量事实账本(ADR-0026)
一次 `AgentRun` 内可能产生**多条**可计费消耗:主模型 completion、外部能力调用
(PDF→MD bundle、音视频转写等)、或代理网关侧旁路调用。这些消耗**不**是子 Run——
它们不持项目锁、不占 session 语义、不复用 `AgentRun` 状态机。它们是 Run 内的
append-only **计量事实**(`UsageFact`)。
`AgentRun` 上的 cost/token 标量是 facts 的派生 rollup cache,不是真相来源;真相在
`UsageFact` 行。这一层抽象让"主模型"与"外部能力"两类消耗共享同一账本,而非给后者
各开一种特化字段或 nested Run。
核心不变式(`PINNED`, ADR-0022 + ADR-0026):
- **append-only** —— fact 一旦写入只读不更不删;`costUsd` 修正走新 fact,不改旧行。
- **`costUsd = none` 表示"未知",不是"零成本"** —— ADR-0022:missing provider cost
remains unknown rather than zero。聚合时不得把 none 当 0 求和。
- **fact 不持锁、不跨 run** —— 它是 Run 内的副作用账本,不是又一次 `@bot` 生命周期。
外部能力调用亦同:它的语义边界是 `capabilityId` + `correlationId`,不是 `RunId`。
- **价目回算是派生** —— `costSource = pricebookDerived` 时,`costUsd` 由
`occurredAt × (provider, model)` 查价目表派生;provider 实报(`providerReported`)
优先于派生。无价目可查时 `costSource = unknown`,`costUsd = none`。
外部能力(capability)注册表、价目表(pricebook)、商业结算均 `OPEN`,pilot 不做
(ADR-0021 仍 defer commercial billing)。本模块只 pin **账本结构与不变式**。
-/
namespace Spec.System.Agent
variable (I : Identifiers) (Time Cost Quantity : Type)
/-- 计量事实类型(`PINNED`, ADR-0026)。集合完整性 `OPEN`——新增 kind 须 surface,
不得静默把外部消耗塞进既有 kind。 -/
inductive UsageFactKind where
/-- 主 agent loop 的模型 completion(`PINNED`)。 -/
| modelCompletion
/-- 外部能力调用(`PINNED`;如 pdf_to_md_bundle、audio_video_to_text)。 -/
| externalCapability
/-- 代理网关侧旁路调用(`PINNED`;非主 loop 的模型调用,如嵌入工具内的子请求)。 -/
| toolProxy
/-- 成本来源(`PINNED`, ADR-0026)。优先级:provider 实报 > 价目派生 > unknown。 -/
inductive CostSource where
/-- provider/gateway 实报(`PINNED`)。 -/
| providerReported
/-- 用 `occurredAt × (provider, model)` 价目表派生(`PINNED`)。 -/
| pricebookDerived
/-- 缺失,而非零(`PINNED`;ADR-0022 missing cost ≠ zero)。 -/
| unknown
/-- 一次可计费消耗的计量事实(`PINNED`, ADR-0026)。append-only;一次 run ≥0 条。
字段分两组:**归属与时间**(run/occurredAt/kind)用于聚合与价目回算;
**计量与成本**(tokens/quantity/unit/costUsd/costSource)用于账目本身。
非 token 计量(页//次)走 `quantity + unit`,与 token 并存而非取代。 -/
structure UsageFact where
/-- 所属 run(`PINNED`;fact 不跨 run,fact 不持锁,ADR-0026)。 -/
run : I.RunId
/-- 消耗发生时刻(`PINNED`);价目回算依赖此字段而非 run.finishedAt,
因为外部能力可能早于 run 结束。 -/
occurredAt : Time
/-- 事实类型(`PINNED`, ADR-0026)。 -/
kind : UsageFactKind
/-- provider 标识(`PINNED`;如 openrouter、mineru、openai_whisper)。 -/
provider : String
/-- 模型标识(`OPEN`;非模型计费可为 none)。 -/
model : Option String
/-- 输入 tokens(`OPEN`;非 token 计量可为 none)。 -/
inputTokens : Option Nat
/-- 输出 tokens(`OPEN`)。 -/
outputTokens : Option Nat
/-- 非 token 计量数值(`OPEN`;如页数、音频秒数)。与 tokens 并存。 -/
quantity : Option Quantity
/-- 计量单位(`OPEN`;如 "pages"、"audio_seconds"、"invocations")。 -/
unit : Option String
/-- 成本(`PINNED`;none = 未知,非零,ADR-0022)。 -/
costUsd : Option Cost
/-- 成本来源(`PINNED`;costUsd ≠ none 时必填,costUsd = none 时为 `unknown`)。 -/
costSource : CostSource
/-- 外部能力标识(`OPEN`;kind = externalCapability 时填,如 pdf_to_md_bundle)。 -/
capabilityId : Option String
/-- 外部请求关联 id(`OPEN`;对账/幂等用,不参与聚合)。 -/
correlationId : Option String
/-- Fact 的成本是否**已知**(`PINNED`, ADR-0026)。聚合 rollup 时,只对已知成本求和;
未知成本的 fact 不贡献 0,而是让汇总保持未知(若所有 fact 均未知)或部分未知。 -/
def UsageFact.HasKnownCost (fact : UsageFact I Time Cost Quantity) : Prop :=
fact.costUsd none
end Spec.System.Agent