feat(tool-bailian-kb): 引入服务缓存支持检索服务动态管理

- 引入 yaml 依赖,更新包依赖配置和锁文件
- 新增 ServiceCache 类,实现检索服务缓存机制,缓存位于本地缓存目录
- 实现服务缓存的读取、写入和异步刷新,支持数据的版本校验和过期控制
- 在插件主入口集成服务缓存,支持按 workspace 维度缓存隔离
- 改写默认 agent_id 获取逻辑,支持单服务自动选取和缓存中的默认服务回退
- 支持检测缓存失效时自动刷新,并能在错误消息中附加当前有效服务列表
- 使用 agent/pre-step 上下文消息注入服务目录,替代工具描述中的服务信息,支持动态更新
- 文档更新完善,说明服务清单的行为语义和管理面与检索面的职责分离
- 为管理面技能引入最佳实践,强调服务命名的重要性和描述字段的价值
- 取消暴露服务清单为模型工具,避免重新引入“先列再搜”的请求环节
This commit is contained in:
zeyu.fz
2026-08-23 14:25:41 +08:00
parent 5f44cb4402
commit 90dd038a2f
20 changed files with 1775 additions and 56 deletions
+46 -3
View File
@@ -76,20 +76,63 @@ Config 同时注册为 `bailian-kb` settings namespace(`installSettingsSection
| `kb_search` | `query`、`agent_id`(**必填**;程序化省略时回退 defaultRetrieveAgentId)、`top_k?`(默认 5,**客户端截断**——服务端无此参数)、`images?` | chunks(text/score/来源)+ total |
| `kb_chat` | `message`、`agent_id`(**必填**;程序化省略时回退 defaultChatAgentId) | 完整答案(内部消费 SSE 流缓冲返回)+ request_id |
服务发现(`kb_service_list` 已移除):通过 `bl knowledge service list` CLI 命令查询可用检索/对话服务及其 agent_id。
两个工具的 **description 保持静态**(不含任何服务 id);可用服务清单由下述服务缓存经 `agent/pre-step` 注入为上下文消息。
## 检索服务缓存与上下文注入
模型要判断"该不该检索",靠的是看到本 workspace 部署了哪些检索服务。插件内部经 `/api/v1/indices/rag/app/list` 拉取该清单并缓存,**不对模型暴露服务发现工具**(`kb_service_list` 不会回归:它会把"先 list 再 search"的额外一轮重新引入);管理面仍用 bl。
### 载体:上下文消息,不是工具描述
清单经 `agent/pre-step` 注入为一条带 source 的 `UserMessage`(`{ kind: 'plugin', plugin: 'tool-bailian-kb/services', form: 'catalog' }`),而不是烘进 tool description。两个原因:
1. **插件加载是每进程一次,不是每会话一次。** 描述在 `apply()` 时定型,长驻宿主里 TTL 只会被评估一次,用户在控制台新建的服务要等重启才能被感知;
2. **重注册工具会废掉 prompt 前缀缓存**(从第一个变化的 schema token 起)。走上下文消息则让 schema 永久稳定。
**变化抑制是正确性要求,不是优化**:`pre-step` 每个“步”(= 一次模型请求)触发一次,一轮里调 5 次工具就触发 6 次。只有清单内容变化时才重发,且判定叠加**可见性**(`session.surface.nodes`)——压缩把清单消息裁掉后会自动重新注入,否则模型会静默失去清单。
### 清单内容策略
| 情形 | 注入内容 |
|---|---|
| 配了默认服务 | 只列该服务 + "另有 N 个" 提示 |
| 未配默认,deployed ≤ 10 | 全量 `agent_id` + 名称 |
| 未配默认,deployed > 10 | 按 `modify_time` 倒序取 10 条,**显式标明截断**与总数 |
| 0 个 / 拉取失败 / 无缓存 | 不注入(工具仍可用) |
英文框架 + 服务名原样保留;空 scene 整节省略;截断必须告知(静默截断会让模型把清单当全集,进而断言"没有对应知识库")。
### 缓存与刷新
落点:`${DSH_HOME:-~/.dsh}/cache/bailian-kb/services-<workspaceId>.json`(临时文件 + `rename()` 原子发布,目录 `0o700`)。按 workspace 分文件是必需的:api key 只能访问自己的 workspace(交叉组合返回 `Endpoint.AccessDenied`),而"自动获取"按钮就是为了切账号。
存:`agent_id` / `agent_name` / `scene` / `status` / `modify_time`,预留 `description`(待后端补齐)。**不存 `pipeline_list`**——实测它常缺 `pipeline_name`、有时整个为空,做不了知识库标签。
| 刷新触发点 | 模型何时看见 |
|---|---|
| pre-step 间隔调度(超 TTL 30 分钟,后台异步,**不阻塞**) | 下一步 |
| 控制台登录成功(`/bailian-kb/autofill` 回调) | 下一步 |
| 调用撞 4xx(agent_id 已失效) | **本步**,刷新后的列表追加进错误消息 |
| workspaceId / apiKey 变更 | 下一步 |
刷新失败只 warn,保留旧文档;并发刷新共享一个请求(pre-step 每步都会检查)。pre-step 监听器**永不抛异常**——抛出会使用户当前这一步失败。未组合 `agents` 的 headless 装配只是没有清单,工具照常可用。
## 错误语义
- HTTP 错误:原始错误透传,模型可通过 `bl knowledge service list` 发现可用服务以纠正无效 `agent_id`;
- HTTP 错误:4xx 时刷新服务缓存并把当前可用服务追加进错误消息(这两个接口上 `agent_id` 是唯一的调用方标识符,所以 4xx 大多是 id 已失效);5xx 与刷新本身失败则原错误透传;
- 凭证缺失:指向 `~/.dsh/.env` / `.credentials.yaml` 配置方式与控制台取 key 页面;
- chat 超时:说明服务端多轮检索特性,建议重试或改用 `kb_search`;
- 服务端错误体截断至 500 字符进入错误信息(优先 `code: message`)。
## 管理面 skill
`skills/bailian-kb/SKILL.md` 随包分发,插件通过 `ctx.inject(['skills'])` 在 skills 服务可用时以 `source: 'bundled'` 运行时注册;无 skills 服务的组合(headless 最小装配)不受影响。内容:bl CLI 安装/鉴权/workspace 解析、建库→上传→部署工作流、agent_id 固定最佳实践。
`skills/bailian-kb/SKILL.md` 随包分发,插件通过 `ctx.inject(['skills'])` 在 skills 服务可用时以 `source: 'bundled'` 运行时注册;无 skills 服务的组合(headless 最小装配)不受影响。文件的 YAML frontmatter 是 name / description 的**单一事实源**,注册时会被剥离(`SkillDefinition.content` 契约上是已去元数据的正文,而 runtime 注册路径不做任何解析)。
内容:bl CLI 安装/鉴权/workspace 解析、建库→上传→部署工作流、服务清单的行为语义。**skill 不承担"该不该检索"的引导**(那是工具描述与上下文清单的事:skill 正文要模型先决定加载才能读到,是二阶决策);它反过来承担一件工具做不到的事:**引导 agent 在 `service create` 时把服务名写清楚**。无 desc 时服务名是唯一语义来源,管理面的动作直接决定检索面的效果。
## Known Limitations
- kb_chat 执行期无进展显示(缓冲式;进展会话事件设计见仓库根 README 与 spec 附录 A)。
- `top_k` 是客户端截断:请求体不含该参数,服务端返回条数由检索服务配置决定,截断只影响进入模型上下文的量。
- **服务画像的质量上限取决于服务名**:`service list` 接口当前不返回描述(已对两个 workspace 实测确认),所以模型只能靠 `agent_name` 判断一个服务能查什么。名字形如 `test-0819` 的部署,引导能力接近于零。后端补齐描述字段后只需改三处(`api-types` 补字段名 → `services.ts` 解析 → `buildServiceCatalog` 追加并截断到 200 字符),缓存已预留 `description` 键,无需迁移。
- 拉取每个 scene 最多 2 页 / 200 条(`page_size` 服务端硬顶 100),超出时标 `truncated` 并在清单里告知。
+12 -2
View File
@@ -1,6 +1,6 @@
{
"name": "@ali/bailian-kb-dsh",
"version": "0.1.15",
"version": "0.1.16",
"description": "Bailian knowledge-base tools for DeepSeek Harness: kb_search and kb_chat over the DashScope RAG API, plus the bl CLI management skill.",
"type": "module",
"main": "lib/index.js",
@@ -38,8 +38,14 @@
"scripts": {
"build": "tsc -b && tsdown"
},
"dependencies": {
"yaml": "^2.4.2"
},
"peerDependencies": {
"@deepseek-ai/cordis": "^4.0.1",
"@deepseek-ai/dsh-agent": "^0.1.0-rc.6",
"@deepseek-ai/dsh-llm": "^0.1.0-rc.6",
"@deepseek-ai/dsh-session": "^0.1.0-rc.6",
"@deepseek-ai/dsh-tools": "^0.1.0-rc.6",
"@deepseek-ai/dsh-credentials": "^0.1.0-rc.6",
"@deepseek-ai/dsh-settings": "^0.1.0-rc.6",
@@ -55,6 +61,9 @@
},
"devDependencies": {
"@deepseek-ai/cordis": "^4.0.1",
"@deepseek-ai/dsh-agent": "^0.1.0-rc.6",
"@deepseek-ai/dsh-llm": "^0.1.0-rc.6",
"@deepseek-ai/dsh-session": "^0.1.0-rc.6",
"@deepseek-ai/dsh-tools": "^0.1.0-rc.6",
"@deepseek-ai/dsh-credentials": "^0.1.0-rc.6",
"@deepseek-ai/dsh-settings": "^0.1.0-rc.6",
@@ -70,6 +79,7 @@
"@types/react": "~18.3.1",
"lightningcss": "^1.32.0",
"react": "^18.2.0",
"tsdown": "^0.22.2"
"tsdown": "^0.22.2",
"yaml": "^2.4.2"
}
}
@@ -6,6 +6,8 @@ description: >-
增删改查 Chunk、管理数据中心类目/文件/集合时使用本 skill。
检索与问答不走本 skill——用原生工具 kb_search(取证据)/ kb_chat(成品问答);
bl knowledge search / chat 仅用于部署后的验证调试(如 --agent-version beta 调试草稿版)。
kb_search / kb_chat 的凭据与工作空间由插件自动解析(~/.dsh/settings.yaml 的 bailian-kb 段、
~/.dsh/.credentials.yaml 的 DASHSCOPE_API_KEY),不要自己去读或传。
普通问答、编程、写作、翻译、泛搜索不触发本 skill。
---
@@ -13,6 +15,15 @@ description: >-
检索面与管理面的分工:**查知识用 `kb_search`(取证据)/ `kb_chat`(成品问答)原生工具;本 skill 只覆盖管理长尾**——知识库全生命周期、文档、检索服务、Chunk、数据中心。
本 skill **不负责**判断何时该检索。可用检索服务的清单(含 agent_id)由插件自动注入到会话上下文里,`kb_search` / `kb_chat` 直接取用;不需要为了检索先加载本 skill。
## 检索服务清单的行为语义
- 清单由插件从百炼 API 拉取后缓存,按会话周期性刷新(约 30 分钟),只含 **deployed** 状态的服务;
- **刚用 `bl` 新建或部署的服务不会立刻出现在清单里**。不用等刷新——命令输出里刚拿到的 `agent_id` 直接可用;
- 服务很多时清单只列最近修改的若干条并标明总数。要找特定服务用 `bl knowledge service list --scene search --name <关键词>`;
- 清单里确实没有能回答用户问题的服务时,如实告知用户,**不要挑一个最像的 agent_id 去试**。
## 前置检查
1. 安装校验:运行 `bl knowledge list --help`。若报 `Unknown command` 或 bl 未安装,执行
@@ -83,8 +94,8 @@ bl knowledge service list --scene search --status deployed # 5.
## 最佳实践
- 用户反复使用同一检索服务时,建议其把 agent_id 写入项目指令(如 AGENTS.md)或让 agent 记住,后续 kb_search / kb_chat 直接携带。
- 服务有 draft/deployed 两种状态:只有 deployed 可被默认版本调用;draft 调试用 `--agent-version beta`。改已发布版本的配置:先改 beta 草稿(`service update`),验证后 `service deploy` 发新版本。
- **建服务时必须把名字写清楚**:`service create --name` 的名称是模型判断"这个服务能查什么"的主要依据(服务描述暂未随列表接口返回)。`检索服务1`、`test-0819` 这类名字会让后续检索无法路由;写成 `产品文档检索`、`HR制度问答` 这种能看出覆盖内容的名字。同时填 `--description`(≤1000 字符),后端补齐列表字段后即可自动生效。
- 服务有 draft/deployed 两种状态:只有 deployed 可被默认版本调用,也只有 deployed 会进入模型看到的服务清单;draft 调试用 `--agent-version beta`。改已发布版本的配置:先改 beta 草稿(`service update`),验证后 `service deploy` 发新版本。
- 导入类命令(`knowledge create`、`doc upload --index-id`、`doc status`)优先带 `--wait` 轮询到终态,避免手工轮询;文档解析失败(如 PARSE_FAILED)会以非零退出码透传错误。
- `chunk add` 有 10 QPS 限流,批量脚本注意节流;响应不带 chunk id,需要 `chunk list` 反查。
- `service list` 必须带 `--scene chat|search`,两个场景要分别查询。
@@ -57,9 +57,9 @@ Usage: bl knowledge service create --name <text> --scene <chat|search> [flags]
| Flag | 说明 |
| --- | --- |
| `--name <text>` | 服务名(≤200 字符,同场景内唯一) |
| `--name <text>` | 服务名(≤200 字符,同场景内唯一)。**名称是模型判断该服务能查什么的主要依据**(描述暂未随 `service list` 返回),写清楚覆盖内容,避开 `检索服务1` / `test-xxx` |
| `--scene <scene>` | chat(问答)或 search(检索) |
| `--description <text>` | 描述(≤1000 字符) |
| `--description <text>` | 描述(≤1000 字符)。建议始终填写:说明覆盖什么内容、适合回答什么问题 |
| `--index-id <id>` | 绑定知识库;其余配置用服务端默认值 |
Notes:
+45
View File
@@ -1,5 +1,50 @@
/** Request/response fields of the DashScope search and chat endpoints, mirrored from the verified bl CLI types. */
/** Retrieval-service scenes; the server requires one per list query. */
export type ServiceScene = 'chat' | 'search'
export interface ServiceListRequest {
agent_scene: ServiceScene
/**
* Verified to be honored by the server, and to mean "deployed or edited"
* (matching the CLI's documented `--status deployed(含 edited)`). The
* spellings `status` and `agent_status_list` are silently ignored.
*/
agent_status?: 'deployed'
agent_name?: string
page_number: number
page_size: number
}
/**
* One row of the service list, mirrored from a verified live response
* (2026-08-23). The response carries NO description field — confirmed against
* two workspaces, including a name-filtered single-row query — even though
* `service create --description` accepts one. `description` is therefore absent
* here until the backend adds it; when it does, **name the field from the real
* response** rather than guessing.
*
* `pipeline_list` is typed but deliberately never consumed: it frequently omits
* `pipeline_name` (leaving an opaque id) and is sometimes empty outright, so it
* cannot serve as a knowledge-base label.
*/
export interface ServiceListRow {
agent_id?: string
agent_name?: string
agent_scene?: string
agent_status?: string
agent_version?: string
create_time?: string
modify_time?: string
pipeline_list?: { pipeline_id?: string; pipeline_name?: string }[]
}
export interface ServiceListResponse {
code?: string
message?: string
data?: { total_count?: number; rows?: ServiceListRow[] }
}
export interface SearchRequest {
query: string
agent_id: string
+7 -1
View File
@@ -1,7 +1,13 @@
/** Protocol path constants and the workspace-subdomain URL builder (external API spec; not configurable). */
/** DashScope knowledge API paths for the two native tools (search, chat); service listing is handled by the bl CLI. */
/**
* DashScope knowledge API paths. `serviceList` backs the plugin's internal
* service cache only — it is deliberately NOT exposed as a model tool (a
* discovery tool reintroduces the "list before you search" round trip this
* design exists to remove); the management surface uses the bl CLI instead.
*/
export const KB_PATHS = {
serviceList: '/api/v1/indices/rag/app/list',
search: '/api/v1/indices/knowledge/search',
chat: '/api/v2/apps/knowledge/chat',
} as const
+88 -11
View File
@@ -13,6 +13,10 @@ import { readBlCliConfig } from './bl-cli.js'
import { consoleLoginState, startConsoleLogin } from './console-login.js'
import { KbClient } from './client.js'
import { registerSkill } from './skill.js'
import { ServiceCache } from './service-cache.js'
import { CATALOG_ENTRY_LIMIT } from './service-catalog.js'
import { installServiceContext } from './service-context.js'
import type { ServiceScene } from './api-types.js'
import { createKbTools } from './tools.js'
/** Minimal webServer route shape (declared inline to avoid a host-package dependency). */
@@ -250,19 +254,69 @@ export function apply(ctx: Context, config: Config): void {
return resolved.value
},
})
/** The workspace id if configured, without the client's guidance throw. */
const resolveWorkspaceIdOrUndefined = async (): Promise<string | undefined> => {
const pinned = current().workspaceId
if (pinned) return pinned
const resolved = await ctx.credentials.resolve(credentialRef('BAILIAN_WORKSPACE_ID'))
return resolved?.value
}
const serviceCache = new ServiceCache({
client,
resolveWorkspaceId: async () => {
const workspaceId = await resolveWorkspaceIdOrUndefined()
if (workspaceId === undefined) throw new Error('workspace id is not configured')
return workspaceId
},
get endpointHost() { return current().endpointHost },
warn: message => { ctx.logger.warn(message) },
})
/** The user's explicitly configured default for one scene: settings layer, then credential. */
const configuredDefaultAgentId = async (scene: ServiceScene): Promise<string | undefined> => {
const pinned = scene === 'search' ? current().defaultRetrieveAgentId : current().defaultChatAgentId
if (pinned) return pinned
const ref = scene === 'search' ? 'BAILIAN_DEFAULT_RETRIEVE_AGENT_ID' : 'BAILIAN_DEFAULT_CHAT_AGENT_ID'
const resolved = await ctx.credentials.resolve(credentialRef(ref))
return resolved?.value
}
/**
* The default service for one scene, falling back to the sole deployed service
* when the workspace has exactly one. That last layer is the zero-configuration
* path for the common 2C deployment: with one service there is nothing to
* choose, so making the user name it in settings buys nothing.
*/
const resolveDefaultAgentId = async (scene: ServiceScene): Promise<string | undefined> => {
const configured = await configuredDefaultAgentId(scene)
if (configured !== undefined) return configured
const workspaceId = await resolveWorkspaceIdOrUndefined()
if (workspaceId === undefined) return undefined
const forScene = serviceCache.peek(workspaceId)?.entries.filter(entry => entry.scene === scene) ?? []
return forScene.length === 1 ? forScene[0]?.agent_id : undefined
}
for (const tool of createKbTools({
client,
resolveDefaultRetrieveAgentId: async () => {
const pinned = current().defaultRetrieveAgentId
if (pinned) return pinned
const resolved = await ctx.credentials.resolve(credentialRef('BAILIAN_DEFAULT_RETRIEVE_AGENT_ID'))
return resolved?.value
},
resolveDefaultChatAgentId: async () => {
const pinned = current().defaultChatAgentId
if (pinned) return pinned
const resolved = await ctx.credentials.resolve(credentialRef('BAILIAN_DEFAULT_CHAT_AGENT_ID'))
return resolved?.value
resolveDefaultRetrieveAgentId: async () => await resolveDefaultAgentId('search'),
resolveDefaultChatAgentId: async () => await resolveDefaultAgentId('chat'),
// Self-heal for a cached id the server has since rejected: refresh once and
// put the current list in the error, which reaches the model this step.
describeServicesAfterRefresh: async (scene) => {
await serviceCache.refresh()
const workspaceId = await resolveWorkspaceIdOrUndefined()
if (workspaceId === undefined) return undefined
const forScene = serviceCache.peek(workspaceId)?.entries.filter(entry => entry.scene === scene) ?? []
if (forScene.length === 0) return undefined
const lines = forScene.slice(0, CATALOG_ENTRY_LIMIT)
.map(entry => `- ${entry.agent_id} — ${entry.agent_name === '' ? '(unnamed)' : entry.agent_name}`)
const more = forScene.length - lines.length
return [
`Deployed ${scene} services in this workspace, re-read just now:`,
...lines,
...(more > 0 ? [`(and ${more} more — \`bl knowledge service list --scene ${scene}\`)`] : []),
].join('\n')
},
get chatTimeoutMs() { return current().chatTimeoutMs },
})) {
@@ -270,6 +324,21 @@ export function apply(ctx: Context, config: Config): void {
}
registerSkill(ctx)
// The service catalog rides an `agent/pre-step` context message rather than the
// tool descriptions: descriptions freeze at plugin load, and a plugin loads
// once per process, so in a long-running host a service created elsewhere
// would never be seen. Optional inject — a headless assembly without `agents`
// simply gets no catalog, and both tools keep working.
ctx.inject(['agents'], (actx) => {
installServiceContext(actx, {
cache: serviceCache,
resolveWorkspaceId: resolveWorkspaceIdOrUndefined,
resolveDefaultRetrieveAgentId: async () => await configuredDefaultAgentId('search'),
resolveDefaultChatAgentId: async () => await configuredDefaultAgentId('chat'),
warn: message => { actx.logger.warn(message) },
})
})
// Export the resolved workspace id as a shell environment variable so
// management CLI commands (`bl knowledge list`, `bl knowledge service list`, etc.)
// running in bash can see the value the settings service resolved.
@@ -393,6 +462,14 @@ export function apply(ctx: Context, config: Config): void {
written.push('workspaceId')
}
if (written.length > 0) await markSeeded(written)
// A completed console login is the one unambiguous signal that the
// account may have changed. Without this the next session would
// build its catalog from the previous account's services, which is
// worse than having no cache at all.
if (written.length > 0) {
serviceCache.invalidate()
void serviceCache.refresh()
}
return written
},
})
@@ -0,0 +1,219 @@
/**
* On-disk cache of the workspace's deployed retrieval services.
*
* Landing spot: `${DSH_HOME:-~/.dsh}/cache/bailian-kb/services-<workspaceId>.json`.
* Per-workspace files are required, not cosmetic: an api key only reaches its own
* workspace (cross-pairs return `Endpoint.AccessDenied`), and the panel's
* "autofill" button exists to switch accounts, so one shared file would blend
* services from different accounts.
*
* Why not `ctx.storage`: the storage hub is absent from every shipped agent
* preset, so `inject(['storage'])` may never fire for a third-party plugin, and
* the JSON backend's on-disk location is decided by its own `root` config — the
* plugin could not tell anyone where the data went. Why not `settings.yaml`: that
* document is the user's, and it hot-reloads, so writing machine-refreshed data
* there both fights the user for the file and republishes configuration for no
* reason.
*
* Every read is best-effort and total: a missing, malformed, foreign, or
* stale-schema file reads as a miss. Callers run inside `agent/pre-step`, where a
* throw fails the user's step.
*/
import { mkdirSync, readFileSync, renameSync, writeFileSync } from 'node:fs'
import { homedir } from 'node:os'
import { dirname, join } from 'node:path'
import type { KbClient } from './client.js'
import { listServices, type ServiceEntry } from './services.js'
/** Bumped whenever the stored shape changes; a mismatch reads as a miss (no migration). */
const CACHE_VERSION = 1
/** Refresh interval. Evaluated per `agent/pre-step`, so a short window genuinely takes effect. */
export const CACHE_TTL_MS = 30 * 60 * 1000
/** The stored document. */
export interface ServiceCacheDocument {
version: number
/** Epoch millis of the fetch that produced `entries`. */
fetchedAt: number
/** Guards against reading a file written for another account or region. */
workspaceId: string
endpointHost: string
entries: ServiceEntry[]
total: number
truncated: boolean
}
/** Resolve the harness home the same way `settings-file` does. */
function dshHome(): string {
const fromEnv = process.env.DSH_HOME
return fromEnv !== undefined && fromEnv !== '' ? fromEnv : join(homedir(), '.dsh')
}
/**
* Build the cache path for one workspace.
* @param workspaceId - the resolved Bailian workspace id.
* @param home - override for tests; defaults to `$DSH_HOME` or `~/.dsh`.
* @returns the absolute file path.
*/
export function serviceCachePath(workspaceId: string, home: string = dshHome()): string {
return join(home, 'cache', 'bailian-kb', `services-${workspaceId}.json`)
}
/**
* Read a cache document, validating it belongs to this workspace and schema.
* @param path - the cache file path.
* @param workspaceId - the workspace the caller is serving.
* @param endpointHost - the host the caller is serving.
* @returns the document, or undefined for any miss (absent, malformed, foreign, or wrong version).
*/
export function readServiceCache(
path: string,
workspaceId: string,
endpointHost: string,
): ServiceCacheDocument | undefined {
let parsed: unknown
try {
parsed = JSON.parse(readFileSync(path, 'utf8'))
} catch (_unreadableOrMalformed) {
return undefined
}
if (typeof parsed !== 'object' || parsed === null || Array.isArray(parsed)) return undefined
const doc = parsed as Partial<ServiceCacheDocument>
if (doc.version !== CACHE_VERSION) return undefined
if (doc.workspaceId !== workspaceId || doc.endpointHost !== endpointHost) return undefined
if (typeof doc.fetchedAt !== 'number' || !Array.isArray(doc.entries)) return undefined
return {
version: CACHE_VERSION,
fetchedAt: doc.fetchedAt,
workspaceId,
endpointHost,
entries: doc.entries,
total: typeof doc.total === 'number' ? doc.total : doc.entries.length,
truncated: doc.truncated === true,
}
}
/**
* Publish a cache document atomically: a reader either sees the previous file or
* the complete new one, never a half-written mix.
* @param path - the cache file path.
* @param doc - the document to store.
*/
export function writeServiceCache(path: string, doc: ServiceCacheDocument): void {
const dir = dirname(path)
mkdirSync(dir, { recursive: true, mode: 0o700 })
const temp = `${path}.${process.pid}.tmp`
writeFileSync(temp, `${JSON.stringify(doc, undefined, 2)}\n`, { mode: 0o600 })
renameSync(temp, path)
}
export interface ServiceCacheOptions {
client: KbClient
/** Resolves the current workspace id; a failure means "not configured yet". */
resolveWorkspaceId: () => Promise<string>
endpointHost: string
/** Reports refresh failures without escalating them. */
warn: (message: string) => void
/** Test seams. */
home?: string
now?: () => number
}
/**
* The service cache: a synchronous in-memory view over the on-disk document,
* plus a deduplicated background refresh.
*/
export class ServiceCache {
/** Last document read or written; undefined until one is available. */
private document: ServiceCacheDocument | undefined
/** The in-flight refresh, if any. One per instance: `pre-step` asks on every model request. */
private inFlight: Promise<void> | undefined
/** Workspace of {@link document}, so a workspace switch invalidates in memory too. */
private loadedFor: string | undefined
constructor(private readonly opts: ServiceCacheOptions) {}
private get now(): number {
return (this.opts.now ?? Date.now)()
}
/**
* The cached entries for a workspace, loading the file on first use.
* Synchronous and total — safe to call from `agent/pre-step`.
* @param workspaceId - the workspace being served.
* @returns the document, or undefined when nothing usable is cached.
*/
peek(workspaceId: string): ServiceCacheDocument | undefined {
if (this.loadedFor !== workspaceId) {
this.document = readServiceCache(
serviceCachePath(workspaceId, this.opts.home ?? dshHome()),
workspaceId,
this.opts.endpointHost,
)
this.loadedFor = workspaceId
}
return this.document
}
/**
* Whether the cached document is missing or older than the TTL.
* @param workspaceId - the workspace being served.
* @returns true when a refresh is due.
*/
isStale(workspaceId: string): boolean {
const doc = this.peek(workspaceId)
return doc === undefined || this.now - doc.fetchedAt >= CACHE_TTL_MS
}
/** Drop the in-memory view and force the next `peek` to re-read from disk. */
invalidate(): void {
this.document = undefined
this.loadedFor = undefined
}
/**
* Fetch and store the current service list.
* Never rejects: failures are warned and leave the previous document in place.
* Concurrent calls share one request.
* @returns a promise resolving once the attempt finishes.
*/
async refresh(): Promise<void> {
// Without this guard `pre-step` would start a fetch on every model request
// while the first is still outstanding.
this.inFlight ??= this.runRefresh().finally(() => {
this.inFlight = undefined
})
return await this.inFlight
}
private async runRefresh(): Promise<void> {
try {
const workspaceId = await this.opts.resolveWorkspaceId()
const list = await listServices(this.opts.client)
// Both scenes failing means the fetch produced nothing; keep the old file.
if (list.failedScenes.length === 2) {
this.opts.warn('bailian-kb service cache not refreshed: both scene queries failed')
return
}
if (list.failedScenes.length > 0) {
this.opts.warn(`bailian-kb service cache refreshed without scene(s): ${list.failedScenes.join(', ')}`)
}
const doc: ServiceCacheDocument = {
version: CACHE_VERSION,
fetchedAt: this.now,
workspaceId,
endpointHost: this.opts.endpointHost,
entries: list.entries,
total: list.total,
truncated: list.truncated,
}
writeServiceCache(serviceCachePath(workspaceId, this.opts.home ?? dshHome()), doc)
this.document = doc
this.loadedFor = workspaceId
} catch (failed) {
this.opts.warn(`bailian-kb service cache refresh failed: ${String(failed)}`)
}
}
}
@@ -0,0 +1,142 @@
/**
* Renders the deployed-service catalog the model reads before deciding whether
* to retrieve.
*
* This is a pure function on purpose: the four selection branches below are the
* whole routing policy, and they are far easier to pin down in tests than
* through a live pre-step.
*
* Rendering conventions (settled):
* - English frame, service names verbatim — same language as the tool
* descriptions, so the model is not switched between languages mid-prompt.
* - Truncation is always stated. Silently cutting the list makes the model treat
* it as complete and flatly answer "there is no such knowledge base".
* - An empty scene omits its whole section. `no chat services` is pure noise and
* invites the model to handle a case that does not exist.
*/
import type { ServiceScene } from './api-types.js'
import type { ServiceEntry } from './services.js'
/** Entries rendered per scene before switching to "most recently modified" mode. */
export const CATALOG_ENTRY_LIMIT = 10
/** Truncation applied to a service description once the backend returns one. */
const DESCRIPTION_LIMIT = 200
export interface CatalogInput {
entries: readonly ServiceEntry[]
/** Server-reported total, which may exceed `entries` when the fetch itself was capped. */
total: number
/** True when the fetch stopped before the server ran out of rows. */
truncated: boolean
defaultRetrieveAgentId?: string
defaultChatAgentId?: string
}
const SCENE_LABEL: Record<ServiceScene, string> = {
search: 'kb_search (retrieval)',
chat: 'kb_chat (grounded Q&A)',
}
/** Render one entry as a single line. */
function renderEntry(entry: ServiceEntry): string {
const name = entry.agent_name === '' ? '(unnamed)' : entry.agent_name
const description = entry.description === undefined || entry.description.trim() === ''
? undefined
: entry.description.trim().length > DESCRIPTION_LIMIT
? `${entry.description.trim().slice(0, DESCRIPTION_LIMIT - 1)}…`
: entry.description.trim()
return `- ${entry.agent_id} — ${name}${description === undefined ? '' : `: ${description}`}`
}
/** Most recently modified first; entries without a timestamp sort last. */
function byRecency(left: ServiceEntry, right: ServiceEntry): number {
const l = left.modify_time ?? ''
const r = right.modify_time ?? ''
if (l === r) return 0
if (l === '') return 1
if (r === '') return -1
return l < r ? 1 : -1
}
/**
* Render one scene's section, or undefined when the scene has no services.
* @param entries - all cached entries (any scene).
* @param scene - the scene to render.
* @param defaultAgentId - this scene's configured default service, when set.
* @param truncatedFetch - whether the fetch itself left rows unread.
* @returns the section lines, or undefined to omit the section entirely.
*/
function renderScene(
entries: readonly ServiceEntry[],
scene: ServiceScene,
defaultAgentId: string | undefined,
truncatedFetch: boolean,
): string[] | undefined {
const forScene = entries.filter(entry => entry.scene === scene)
if (forScene.length === 0) return undefined
const lines = [`${SCENE_LABEL[scene]}:`]
const configured = defaultAgentId === undefined
? undefined
: forScene.find(entry => entry.agent_id === defaultAgentId)
if (configured !== undefined) {
// A configured default is the user's own pick: the highest-quality signal
// available, so it is the only entry worth spending context on.
lines.push(renderEntry(configured))
const others = forScene.length - 1
if (others > 0) {
lines.push(
` (default service; ${others} other${others === 1 ? '' : 's'} exist — `
+ `run \`bl knowledge service list --scene ${scene}\` to see them)`,
)
}
return lines
}
if (forScene.length <= CATALOG_ENTRY_LIMIT && !truncatedFetch) {
lines.push(...forScene.map(renderEntry))
return lines
}
const shown = [...forScene].sort(byRecency).slice(0, CATALOG_ENTRY_LIMIT)
lines.push(...shown.map(renderEntry))
// State the shortfall: the model must know this list is partial before it
// concludes no service covers the question.
const knownTotal = Math.max(forScene.length, shown.length)
lines.push(
` (showing ${shown.length} most recently modified of ${truncatedFetch ? 'more than ' : ''}`
+ `${knownTotal} deployed ${scene} services — run \`bl knowledge service list --scene ${scene} `
+ '--name <keyword>\` to look for others)',
)
return lines
}
/**
* Build the catalog text for one cached service list.
* @param input - the cached entries plus the deployment's configured defaults.
* @returns the model-facing text, or undefined when there is nothing worth injecting.
*/
export function buildServiceCatalog(input: CatalogInput): string | undefined {
const search = renderScene(input.entries, 'search', input.defaultRetrieveAgentId, input.truncated)
const chat = renderScene(input.entries, 'chat', input.defaultChatAgentId, input.truncated)
if (search === undefined && chat === undefined) return undefined
return [
'<system-reminder>',
// The header must not name the tools: a scene with no services omits its
// section, and naming that tool anyway would invite passing an id from the
// other scene, which the service rejects. Section labels carry the mapping.
'Bailian knowledge services deployed in this workspace, grouped by the tool that accepts them. '
+ "Pass an id from the matching section as that tool's `agent_id` argument — it is required and "
+ 'cannot be guessed.',
'',
...(search ?? []),
...(search !== undefined && chat !== undefined ? [''] : []),
...(chat ?? []),
'',
'If none of these services covers what the user is asking about, say so plainly rather than '
+ 'trying the closest-looking id — an unrelated retrieval result is worse than none.',
'</system-reminder>',
].join('\n')
}
@@ -0,0 +1,143 @@
/**
* Publishes the cached service catalog into each request as a sourced context
* message, on the `agent/pre-step` waterfall.
*
* Why a context message rather than the tool descriptions: a tool description is
* fixed when the plugin loads, and a plugin loads once per PROCESS, not once per
* session. In a long-running host the TTL would be evaluated exactly once at
* `apply()` and a service created elsewhere would never be noticed until a
* restart. Re-registering tools to refresh a description instead invalidates the
* prompt prefix cache from the first changed schema token. Injecting context
* keeps the tool schemas byte-stable forever and still refreshes per step.
*
* Two hard constraints follow from `agent/pre-step` semantics:
*
* 1. `pre-step` fires once per STEP, and a step is one model request — a turn with
* five tool calls fires it six times. Re-injecting each time would insert six
* copies into one turn and void the KV cache from the first insertion onward,
* so change suppression is a correctness requirement, not an optimization.
* 2. A throwing listener fails the proposed step, i.e. the user's turn stalls.
* Everything here is therefore wrapped: any failure degrades to "inject
* nothing this step".
*/
import type { Context } from '@deepseek-ai/cordis'
import type { Agent, PreStepDecision } from '@deepseek-ai/dsh-agent'
import { createUserMessage } from '@deepseek-ai/dsh-llm'
import type { UserMessage } from '@deepseek-ai/dsh-session'
import { buildServiceCatalog } from './service-catalog.js'
import type { ServiceCache } from './service-cache.js'
/** Marks this plugin's own injections in the durable log. */
const SOURCE_PLUGIN = 'tool-bailian-kb/services'
export interface ServiceContextOptions {
cache: ServiceCache
/** Resolves the workspace being served; undefined means "not configured yet". */
resolveWorkspaceId: () => Promise<string | undefined>
/** Resolves the configured default retrieval service (settings then credential). */
resolveDefaultRetrieveAgentId: () => Promise<string | undefined>
/** Resolves the configured default chat service (settings then credential). */
resolveDefaultChatAgentId: () => Promise<string | undefined>
warn: (message: string) => void
}
/** Whether one durable message came from this module. */
function isOwnInjection(source: { kind: string; plugin?: string }): boolean {
return source.kind === 'plugin' && source.plugin === SOURCE_PLUGIN
}
/** Concatenate a message's text parts, which is what the model actually reads. */
function messageText(message: UserMessage): string {
return message.content
.flatMap(part => part.type === 'text' ? [part.text] : [])
.join('\n')
}
/**
* The catalog text this session last injected AND still shows the model.
*
* The visibility test is the subtle half. Scanning only for "did we ever publish
* this" would make compaction permanent data loss: once the catalog message is
* dropped from the surface, an identical digest would suppress every future
* injection and the model would silently spend the rest of the session without a
* service list.
* @param agent - the subject agent.
* @returns the visible catalog text, or undefined when none is currently visible.
*/
function visibleCatalogText(agent: Agent): string | undefined {
const visible = new Set(agent.session.surface.nodes)
const events = agent.session.events
for (let index = events.length - 1; index >= 0; index -= 1) {
const event = events[index]
if (event === undefined) continue
if (event.type !== 'user/message' || !isOwnInjection(event.data.source)) continue
return visible.has(event.seq) ? messageText(event.data) : undefined
}
return undefined
}
/** This module's proposed-but-not-yet-entered message, if the batch already carries one. */
function pendingCatalog(messages: readonly UserMessage[]): UserMessage | undefined {
return messages.find(message => isOwnInjection(message.source))
}
/**
* Install the pre-step listener that keeps the catalog present and current.
* @param ctx - a context with `agents` available.
* @param opts - the cache plus the deployment's resolved workspace and defaults.
*/
export function installServiceContext(ctx: Context, opts: ServiceContextOptions): void {
ctx.on('agent/pre-step', async ({ agent, signal }, next): Promise<PreStepDecision> => {
const decision = await next()
if (decision.kind === 'reject' || signal.aborted) return decision
try {
const workspaceId = await opts.resolveWorkspaceId()
// Nothing is configured yet: the tools themselves will explain that.
if (workspaceId === undefined || workspaceId === '') return decision
// Refresh scheduling lives here, not at plugin load, so a long-running
// process still notices services created elsewhere. Never awaited: a slow
// list request must not delay the user's request.
if (opts.cache.isStale(workspaceId)) void opts.cache.refresh()
const document = opts.cache.peek(workspaceId)
if (document === undefined) return decision
const [defaultRetrieveAgentId, defaultChatAgentId] = await Promise.all([
opts.resolveDefaultRetrieveAgentId(),
opts.resolveDefaultChatAgentId(),
])
if (signal.aborted) return decision
const text = buildServiceCatalog({
entries: document.entries,
total: document.total,
truncated: document.truncated,
...(defaultRetrieveAgentId !== undefined ? { defaultRetrieveAgentId } : {}),
...(defaultChatAgentId !== undefined ? { defaultChatAgentId } : {}),
})
if (text === undefined) return decision
// Identical to what the model already sees: stay out of the way. This is
// the branch that runs on nearly every step.
if (visibleCatalogText(agent) === text) return decision
const pending = pendingCatalog(decision.messages)
if (pending !== undefined && messageText(pending) === text) return decision
const catalog = createUserMessage({
content: [{ type: 'text', text }],
source: { kind: 'plugin', plugin: SOURCE_PLUGIN, form: 'catalog' },
})
return {
kind: 'enter',
messages: pending === undefined
? [...decision.messages, catalog]
: decision.messages.map(message => message.id === pending.id ? catalog : message),
}
} catch (failed) {
// A throw here would fail the user's step; a missing catalog is far cheaper.
opts.warn(`bailian-kb service catalog not injected: ${String(failed)}`)
return decision
}
}, { prepend: true })
}
+113
View File
@@ -0,0 +1,113 @@
/**
* Retrieval-service discovery for the plugin's internal cache. Not a model tool:
* see `KB_PATHS.serviceList`.
*
* Two verified server facts shape this module:
* - `page_size` is capped at 100 regardless of what is requested, so a workspace
* with hundreds of services needs many round trips. Test/CI workspaces reach
* the high hundreds (913 observed), which is pure noise for routing, so this
* module stops after {@link MAX_PAGES} and reports the shortfall instead of
* faithfully paging through it.
* - `agent_status: 'deployed'` is honored and means "deployed or edited". Only
* those are callable by the default agent version, so drafts never reach the
* model.
*/
import type { ServiceListResponse, ServiceScene } from './api-types.js'
import type { KbClient } from './client.js'
import { KB_PATHS } from './endpoints.js'
/** Server page-size maximum; larger requests are silently clamped to this. */
const PAGE_SIZE = 100
/** Pages fetched per scene before reporting truncation (200 rows is far past the useful range). */
const MAX_PAGES = 2
/** One deployed retrieval or Q&A service, reduced to the fields that inform routing. */
export interface ServiceEntry {
agent_id: string
agent_name: string
scene: ServiceScene
/** `deployed` or `edited` — both are callable by the default version. */
status: string
/** Last modification timestamp; the only signal for "which of these is in use". */
modify_time?: string
/** Absent until the backend adds a description to the list response. */
description?: string
}
export interface ServiceList {
entries: ServiceEntry[]
/** Server-reported total across the queried scenes, including rows never fetched. */
total: number
/** True when a scene reported more rows than {@link MAX_PAGES} pages returned. */
truncated: boolean
/** Scenes whose query failed; a partial list stays usable. */
failedScenes: ServiceScene[]
}
const SCENES: readonly ServiceScene[] = ['search', 'chat']
/**
* Fetch the deployed services of one scene, stopping at the page cap.
* @param client - the shared knowledge API client.
* @param scene - `search` or `chat`.
* @returns the scene's entries, its server-reported total, and whether rows were left unfetched.
*/
async function listScene(
client: KbClient,
scene: ServiceScene,
): Promise<{ entries: ServiceEntry[]; total: number; truncated: boolean }> {
const entries: ServiceEntry[] = []
let total = 0
for (let page = 1; page <= MAX_PAGES; page += 1) {
const res = await client.postJson<ServiceListResponse>(KB_PATHS.serviceList, {
agent_scene: scene,
agent_status: 'deployed',
page_number: page,
page_size: PAGE_SIZE,
})
total = res.data?.total_count ?? total
const rows = res.data?.rows ?? []
for (const row of rows) {
const agentId = row.agent_id ?? ''
// A row without an id cannot be called, so it has no reason to exist here.
if (agentId === '') continue
entries.push({
agent_id: agentId,
agent_name: row.agent_name ?? '',
scene,
status: row.agent_status ?? '',
...(typeof row.modify_time === 'string' ? { modify_time: row.modify_time } : {}),
})
}
// A short page is the last page; the server has nothing further to give.
if (rows.length < PAGE_SIZE) return { entries, total, truncated: false }
}
return { entries, total, truncated: total > entries.length }
}
/**
* List the deployed services of both scenes.
* A scene that fails is recorded and skipped rather than failing the whole
* refresh: half a list still routes better than none.
* @param client - the shared knowledge API client.
* @returns merged entries plus totals, truncation, and per-scene failures.
*/
export async function listServices(client: KbClient): Promise<ServiceList> {
const entries: ServiceEntry[] = []
const failedScenes: ServiceScene[] = []
let total = 0
let truncated = false
for (const scene of SCENES) {
try {
const result = await listScene(client, scene)
entries.push(...result.entries)
total += result.total
truncated = truncated || result.truncated
} catch (_sceneFailed) {
failedScenes.push(scene)
}
}
return { entries, total, truncated, failedScenes }
}
+111 -10
View File
@@ -1,31 +1,132 @@
/** Runtime skill registration: the packaged bl-management SKILL.md joins the catalog when a skills registry is composed. */
/**
* Runtime skill registration: the packaged bailian-kb SKILL.md joins the catalog
* when a skills registry is composed.
*
* The file's YAML frontmatter is the single source of truth for the routing name
* and description — duplicating them here drifts, and the copy that loses is the
* one nobody reads. The frontmatter must also be STRIPPED from the registered
* body: `SkillDefinition.content` is contractually the body a provider has
* already cleaned of its own metadata, and `ctx.skills.register()` parses
* nothing, so handing over the raw file ships the YAML block into the model's
* context.
*/
import { readFileSync } from 'node:fs'
import { join } from 'node:path'
import { fileURLToPath } from 'node:url'
import { parse as parseYaml } from 'yaml'
import type { Context } from '@deepseek-ai/cordis'
// Type-only: resolves ctx.skills for the optional inject below.
import type {} from '@deepseek-ai/dsh-skill'
const SKILL_DIR = fileURLToPath(new URL('../skills/bailian-kb/', import.meta.url))
/** The registrable fields carried by one skill file. */
export interface ParsedSkillFile {
/** Kebab-case skill name from frontmatter. */
name: string
/** Routing description from frontmatter (the catalog truncates at 500 chars). */
description: string
/** Optional extra routing guidance. */
whenToUse?: string
/** Optional frontmatter `metadata` object. */
metadata?: Record<string, unknown>
/** Markdown body with the frontmatter block removed. */
content: string
}
/**
* Locate the frontmatter block, mirroring the filesystem provider's delimiters
* so a file that loads from disk behaves identically when bundled.
* @param raw - the file's full text.
* @returns the frontmatter YAML and the body after it, or undefined when unfenced.
*/
function splitFrontmatter(raw: string): { yaml: string; body: string } | undefined {
const firstLineEnd = raw.indexOf('\n')
if (firstLineEnd < 0) return undefined
if (raw.slice(0, firstLineEnd).replace(/\r$/, '') !== '---') return undefined
const start = firstLineEnd + 1
let lineStart = start
while (lineStart <= raw.length) {
const nextNewline = raw.indexOf('\n', lineStart)
const lineEnd = nextNewline < 0 ? raw.length : nextNewline
if (raw.slice(lineStart, lineEnd).replace(/\r$/, '') === '---') {
return {
yaml: raw.slice(start, lineStart),
body: raw.slice(nextNewline < 0 ? raw.length : nextNewline + 1),
}
}
if (nextNewline < 0) return undefined
lineStart = nextNewline + 1
}
return undefined
}
/**
* Parse one skill file into its registrable fields.
* @param raw - the file's full text.
* @returns the parsed fields, or undefined when the frontmatter is absent, unparsable, or missing name/description.
*/
export function parseSkillFile(raw: string): ParsedSkillFile | undefined {
const split = splitFrontmatter(raw)
if (split === undefined) return undefined
let data: unknown
try {
data = parseYaml(split.yaml)
} catch (_invalidYaml) {
return undefined
}
if (typeof data !== 'object' || data === null || Array.isArray(data)) return undefined
const record = data as Record<string, unknown>
const name = typeof record.name === 'string' ? record.name.trim() : ''
const description = typeof record.description === 'string' ? record.description.trim() : ''
// The registry rejects a blank description outright; failing here keeps the
// diagnostic on the file instead of on the registration call.
if (name === '' || description === '') return undefined
const whenToUse = typeof record.whenToUse === 'string' ? record.whenToUse.trim() : ''
const metadata = typeof record.metadata === 'object' && record.metadata !== null && !Array.isArray(record.metadata)
? record.metadata as Record<string, unknown>
: undefined
return {
name,
description,
...(whenToUse !== '' ? { whenToUse } : {}),
...(metadata !== undefined ? { metadata } : {}),
content: split.body,
}
}
/**
* Register the management skill when the skills registry is composed; headless
* assemblies without the seam stay unaffected.
* assemblies without the seam stay unaffected. An unreadable or malformed file
* degrades to a warning — a broken bundled asset must not fail plugin load.
* @param ctx - the plugin context.
*/
export function registerSkill(ctx: Context): void {
ctx.inject(['skills'], (skillCtx) => {
const content = readFileSync(join(SKILL_DIR, 'SKILL.md'), 'utf8')
const path = join(SKILL_DIR, 'SKILL.md')
let raw: string
try {
raw = readFileSync(path, 'utf8')
} catch (unreadable) {
skillCtx.logger.warn(`bailian-kb skill not registered: cannot read ${path}: ${String(unreadable)}`)
return
}
const parsed = parseSkillFile(raw)
if (parsed === undefined) {
skillCtx.logger.warn(
`bailian-kb skill not registered: ${path} needs YAML frontmatter carrying name and description`,
)
return
}
skillCtx.skills.register({
name: 'bailian-kb',
description:
'Manage Bailian knowledge bases with the bl CLI: create/update KBs, upload documents, deploy '
+ 'retrieval services, and maintain chunks. Retrieval itself uses the native kb_search/kb_chat tools. '
+ 'Credentials and workspace for kb_search/kb_chat resolve automatically from DSH config '
+ '(bailian-kb in ~/.dsh/settings.yaml, DASHSCOPE_API_KEY in ~/.dsh/.credentials.yaml).',
content,
name: parsed.name,
description: parsed.description,
...(parsed.whenToUse !== undefined ? { whenToUse: parsed.whenToUse } : {}),
...(parsed.metadata !== undefined ? { metadata: parsed.metadata } : {}),
content: parsed.content,
source: 'bundled',
path,
resourceBase: { kind: 'directory', path: SKILL_DIR },
})
})
+51 -12
View File
@@ -5,6 +5,13 @@
* fallback to a configured default (settings/config or credential) is retained as defense-in-depth,
* but note defineTool validates args against the schema before execute, so through that entry point
* the fallback is inert; the model-facing contract is explicit.
*
* These descriptions are deliberately STATIC. The available service ids are
* deployment state that changes while the process runs, and re-registering a tool
* to refresh its description invalidates the prompt prefix cache from the first
* changed schema token. The live catalog therefore rides an `agent/pre-step`
* context message instead (see `service-context.ts`), leaving these schemas
* byte-stable for the life of the process.
*/
import { defineTool } from '@deepseek-ai/dsh-tools'
@@ -22,17 +29,41 @@ export interface KbToolDeps {
resolveDefaultRetrieveAgentId?: () => Promise<string | undefined>
/** Resolves the default chat agent id per call (settings/patch config or credential); omitted means no default for kb_chat. */
resolveDefaultChatAgentId?: () => Promise<string | undefined>
/**
* Refreshes the service cache and summarizes what the workspace currently
* deploys for one scene. Called only after a client-side API failure, so a
* stale cached id self-corrects within the same step instead of waiting for the
* next scheduled refresh.
*/
describeServicesAfterRefresh?: (scene: 'search' | 'chat') => Promise<string | undefined>
/** Read per call (a live-settings deployment supplies a getter). */
chatTimeoutMs: number
}
/**
* Forward the original tool error unchanged; service discovery now lives
* in the bl management skill (`bl knowledge service list`), so the tool no
* longer makes a best-effort API round-trip to enrich the message.
* Rethrow an API failure, appending a freshly refreshed service list when the
* server rejected the request.
*
* On these two endpoints `agent_id` is the only caller-supplied identifier, so a
* 4xx is most often a service id that no longer exists — the recovery the model
* needs is the current list, delivered in the error message. The message is
* ordinary conversation text appended at the tail, so unlike a re-registered
* description it does not disturb the request prefix.
*/
async function withServiceHint(_client: KbClient, err: unknown): Promise<never> {
throw err
async function withServiceHint(
err: unknown,
scene: 'search' | 'chat',
describe: KbToolDeps['describeServicesAfterRefresh'],
): Promise<never> {
if (describe === undefined || !(err instanceof KbApiError)) throw err
const status = err.status
if (status === undefined || status < 400 || status >= 500) throw err
// Best-effort enrichment: a failing refresh must not replace the real error.
const summary = await describe(scene).catch(() => undefined)
if (summary === undefined) throw err
throw new KbApiError(`${err.message}
${summary}`, status)
}
/**
@@ -48,9 +79,12 @@ export function createKbTools(deps: KbToolDeps) {
type: 'string' as const,
required: true as const,
description: 'Retrieval/Q&A service id. REQUIRED: the schema cannot know whether this deployment '
+ 'configures a default service, so always pass one. Find ids via '
+ '`bl knowledge service list --scene search --workspace-id <workspaceId>` (workspaceId resolves '
+ 'automatically from DSH settings: bailian-kb.workspaceId in ~/.dsh/settings.yaml).',
+ 'configures a default service, so always pass one. The deployed services of this workspace, '
+ 'with their ids, are listed in a context message in this conversation; take the id from the '
+ 'section matching the tool you are calling. If that list is absent or none of its services '
+ 'covers the question, run `bl knowledge service list --scene search --name <keyword>` to look '
+ '(workspaceId resolves automatically from DSH settings: bailian-kb.workspaceId in '
+ '~/.dsh/settings.yaml).',
}
const resolveRetrieveAgentId = async (supplied: string | undefined): Promise<string> => {
if (supplied !== undefined) return supplied
@@ -89,7 +123,9 @@ export function createKbTools(deps: KbToolDeps) {
+ 'Use kb_chat instead when the user question can be answered by the knowledge base alone. '
+ 'Credentials and workspace resolve automatically from DSH config '
+ '(bailian-kb in ~/.dsh/settings.yaml, DASHSCOPE_API_KEY in ~/.dsh/.credentials.yaml) — '
+ 'never read or pass them yourself. agent_id is REQUIRED (see its parameter description).',
+ 'never read or pass them yourself. agent_id is REQUIRED (see its parameter description). '
+ 'If no listed service covers what the user is asking about, say so plainly rather than trying '
+ 'the closest-looking id: unrelated evidence is worse for the user than none.',
parameters: {
query: { type: 'string', required: true, description: 'Search query text.' },
agent_id: agentIdParam,
@@ -134,7 +170,8 @@ export function createKbTools(deps: KbToolDeps) {
...(client.agentVersion ? { agent_version: client.agentVersion } : {}),
...(args.images && args.images.length > 0 ? { images: args.images } : {}),
}
const res = await client.postJson<SearchResponse>(KB_PATHS.search, body).catch(err => withServiceHint(client, err))
const res = await client.postJson<SearchResponse>(KB_PATHS.search, body)
.catch(async err => await withServiceHint(err, 'search', deps.describeServicesAfterRefresh))
const nodes = (res.data?.nodes ?? []).slice(0, topK)
return {
chunks: nodes.map(n => ({
@@ -160,7 +197,9 @@ export function createKbTools(deps: KbToolDeps) {
+ 'The pipeline runs an internal analysis/retrieval loop and may take a few minutes. '
+ 'Credentials and workspace resolve automatically from DSH config '
+ '(bailian-kb in ~/.dsh/settings.yaml, DASHSCOPE_API_KEY in ~/.dsh/.credentials.yaml) — '
+ 'never read or pass them yourself. agent_id is REQUIRED (see its parameter description).',
+ 'never read or pass them yourself. agent_id is REQUIRED (see its parameter description). '
+ 'If no listed service covers what the user is asking about, say so plainly rather than trying '
+ 'the closest-looking id.',
parameters: {
message: { type: 'string', required: true, description: 'The question to ask.' },
agent_id: agentIdParam,
@@ -196,7 +235,7 @@ export function createKbTools(deps: KbToolDeps) {
+ 'and long questions can exceed the deployment timeout. Retry, or use kb_search for raw chunks instead.',
)
}
return await withServiceHint(client, err)
return await withServiceHint(err, 'chat', deps.describeServicesAfterRefresh)
}
const { answer, requestId } = await consumeChatStream(res)
return { answer, ...(requestId ? { request_id: requestId } : {}) }
@@ -0,0 +1,195 @@
import { mkdtempSync, readFileSync, readdirSync, statSync, writeFileSync, mkdirSync } from 'node:fs'
import { tmpdir } from 'node:os'
import { dirname, join } from 'node:path'
import { describe, expect, it, vi } from 'vitest'
import type { KbClient } from '../src/client.js'
import {
CACHE_TTL_MS,
ServiceCache,
readServiceCache,
serviceCachePath,
writeServiceCache,
type ServiceCacheDocument,
} from '../src/service-cache.js'
const HOST = 'cn-beijing.maas.aliyuncs.com'
function tempHome(): string {
return mkdtempSync(join(tmpdir(), 'bailian-kb-cache-'))
}
function doc(overrides: Partial<ServiceCacheDocument> = {}): ServiceCacheDocument {
return {
version: 1,
fetchedAt: 1_000,
workspaceId: 'llm-a',
endpointHost: HOST,
entries: [{ agent_id: 'aid-1', agent_name: 'svc', scene: 'search', status: 'deployed' }],
total: 1,
truncated: false,
...overrides,
}
}
/** Write raw text to the cache path, bypassing the writer's validation. */
function seedRaw(home: string, workspaceId: string, text: string): string {
const path = serviceCachePath(workspaceId, home)
mkdirSync(dirname(path), { recursive: true })
writeFileSync(path, text)
return path
}
describe('serviceCachePath', () => {
it('separates workspaces by filename', () => {
// An api key only reaches its own workspace, so one shared file would blend accounts.
expect(serviceCachePath('llm-a', '/home')).toBe('/home/cache/bailian-kb/services-llm-a.json')
expect(serviceCachePath('llm-b', '/home')).not.toBe(serviceCachePath('llm-a', '/home'))
})
})
describe('readServiceCache', () => {
it('round-trips a document written by writeServiceCache', () => {
const home = tempHome()
const path = serviceCachePath('llm-a', home)
writeServiceCache(path, doc())
expect(readServiceCache(path, 'llm-a', HOST)).toEqual(doc())
})
it('reads every unusable file as a miss instead of throwing', () => {
const home = tempHome()
// Absent.
expect(readServiceCache(serviceCachePath('nope', home), 'nope', HOST)).toBeUndefined()
// Malformed JSON.
expect(readServiceCache(seedRaw(home, 'a', '{oops'), 'a', HOST)).toBeUndefined()
// Non-object roots.
expect(readServiceCache(seedRaw(home, 'b', '["x"]'), 'b', HOST)).toBeUndefined()
// Newer or older schema: treated as a miss, never migrated.
expect(readServiceCache(seedRaw(home, 'c', JSON.stringify(doc({ version: 2 }))), 'c', HOST)).toBeUndefined()
// Required fields of the wrong type.
expect(readServiceCache(seedRaw(home, 'd', JSON.stringify(doc({ fetchedAt: 'soon' as never }))), 'd', HOST))
.toBeUndefined()
})
it('refuses a document belonging to another workspace or host', () => {
const home = tempHome()
const path = seedRaw(home, 'llm-a', JSON.stringify(doc({ workspaceId: 'llm-other' })))
expect(readServiceCache(path, 'llm-a', HOST)).toBeUndefined()
const hostPath = seedRaw(home, 'llm-b', JSON.stringify(doc({ workspaceId: 'llm-b', endpointHost: 'other.host' })))
expect(readServiceCache(hostPath, 'llm-b', HOST)).toBeUndefined()
})
})
describe('writeServiceCache', () => {
it('publishes atomically and leaves no temp file behind', () => {
const home = tempHome()
const path = serviceCachePath('llm-a', home)
writeServiceCache(path, doc())
const names = readdirSync(dirname(path))
expect(names).toEqual(['services-llm-a.json'])
expect(JSON.parse(readFileSync(path, 'utf8'))).toEqual(doc())
})
it('creates the cache directory owner-only', () => {
const home = tempHome()
const path = serviceCachePath('llm-a', home)
writeServiceCache(path, doc())
expect(statSync(dirname(path)).mode & 0o777).toBe(0o700)
})
})
/** A cache wired to a client returning one page per scene. */
function cacheWith(home: string, onPost: () => Promise<unknown>, now: () => number = () => 5_000) {
const warn = vi.fn()
const postJson = vi.fn(onPost)
const cache = new ServiceCache({
client: { postJson } as unknown as KbClient,
resolveWorkspaceId: async () => 'llm-a',
endpointHost: HOST,
warn,
home,
now,
})
return { cache, warn, postJson }
}
describe('ServiceCache', () => {
const emptyPage = { code: 'Success', data: { total_count: 0, rows: [] } }
it('treats a missing document as stale and a fresh one as current', () => {
const home = tempHome()
const { cache } = cacheWith(home, async () => emptyPage)
expect(cache.isStale('llm-a')).toBe(true)
writeServiceCache(serviceCachePath('llm-a', home), doc({ fetchedAt: 5_000 }))
cache.invalidate()
expect(cache.isStale('llm-a')).toBe(false)
})
it('goes stale once the TTL elapses', () => {
const home = tempHome()
writeServiceCache(serviceCachePath('llm-a', home), doc({ fetchedAt: 0 }))
const { cache } = cacheWith(home, async () => emptyPage, () => CACHE_TTL_MS)
expect(cache.isStale('llm-a')).toBe(true)
})
it('re-reads from disk when the workspace changes', () => {
const home = tempHome()
writeServiceCache(serviceCachePath('llm-a', home), doc({ fetchedAt: 5_000 }))
const { cache } = cacheWith(home, async () => emptyPage)
expect(cache.peek('llm-a')?.workspaceId).toBe('llm-a')
// Switching accounts must not keep serving the previous workspace's list.
expect(cache.peek('llm-b')).toBeUndefined()
})
it('shares one in-flight request across concurrent refreshes', async () => {
const home = tempHome()
// pre-step runs on every model request, so an unguarded refresh would pile up.
const { cache, postJson } = cacheWith(home, async () => emptyPage)
await Promise.all([cache.refresh(), cache.refresh(), cache.refresh()])
// Two calls total: one per scene, from a single shared refresh.
expect(postJson).toHaveBeenCalledTimes(2)
})
it('stores a fetched list and serves it synchronously afterwards', async () => {
const home = tempHome()
const { cache } = cacheWith(home, async () => ({
code: 'Success',
data: { total_count: 1, rows: [{ agent_id: 'aid-9', agent_name: 'svc', agent_status: 'deployed' }] },
}))
await cache.refresh()
const stored = cache.peek('llm-a')
expect(stored?.entries.map(e => e.agent_id)).toEqual(['aid-9', 'aid-9']) // one per scene
expect(stored?.fetchedAt).toBe(5_000)
// And it survives as a file for the next process.
expect(readServiceCache(serviceCachePath('llm-a', home), 'llm-a', HOST)).toEqual(stored)
})
it('never rejects on failure and keeps the previous document', async () => {
const home = tempHome()
writeServiceCache(serviceCachePath('llm-a', home), doc({ fetchedAt: 5_000 }))
const { cache, warn } = cacheWith(home, async () => {
throw new Error('network down')
})
expect(cache.peek('llm-a')?.entries).toHaveLength(1)
await expect(cache.refresh()).resolves.toBeUndefined()
expect(warn).toHaveBeenCalled()
// The stale-but-usable list is still there.
expect(cache.peek('llm-a')?.entries).toHaveLength(1)
})
it('does not reject when the workspace is not configured yet', async () => {
const home = tempHome()
const warn = vi.fn()
const cache = new ServiceCache({
client: { postJson: vi.fn() } as unknown as KbClient,
resolveWorkspaceId: async () => {
throw new Error('workspace id is not configured')
},
endpointHost: HOST,
warn,
home,
})
await expect(cache.refresh()).resolves.toBeUndefined()
expect(warn).toHaveBeenCalledWith(expect.stringContaining('refresh failed'))
})
})
@@ -0,0 +1,115 @@
import { describe, expect, it } from 'vitest'
import { CATALOG_ENTRY_LIMIT, buildServiceCatalog } from '../src/service-catalog.js'
import type { ServiceEntry } from '../src/services.js'
function entry(overrides: Partial<ServiceEntry> & { agent_id: string }): ServiceEntry {
return {
agent_name: `name-${overrides.agent_id}`,
scene: 'search',
status: 'deployed',
...overrides,
}
}
function catalog(entries: ServiceEntry[], extra: Partial<Parameters<typeof buildServiceCatalog>[0]> = {}) {
return buildServiceCatalog({ entries, total: entries.length, truncated: false, ...extra })
}
describe('buildServiceCatalog', () => {
it('returns undefined when there is nothing worth injecting', () => {
// No services at all: the tool descriptions alone keep the tools usable.
expect(catalog([])).toBeUndefined()
})
it('lists every service when the count is within the limit', () => {
const text = catalog([
entry({ agent_id: 'aid-1', agent_name: 'RAG学习-检索' }),
entry({ agent_id: 'aid-2', agent_name: '产品文档' }),
])
expect(text).toContain('aid-1 — RAG学习-检索')
expect(text).toContain('aid-2 — 产品文档')
// Nothing was cut, so nothing should claim otherwise.
expect(text).not.toContain('showing')
})
it('omits a scene with no services entirely', () => {
const text = catalog([entry({ agent_id: 'aid-1' })])
expect(text).toContain('kb_search')
// "no chat services" is noise that invites handling a case that does not exist.
expect(text).not.toContain('kb_chat')
expect(text).not.toMatch(/no chat/i)
})
it('renders both scenes without mixing their services', () => {
const text = catalog([
entry({ agent_id: 'aid-s', scene: 'search' }),
entry({ agent_id: 'aid-c', scene: 'chat' }),
]) ?? ''
const searchAt = text.indexOf('kb_search')
const chatAt = text.indexOf('kb_chat')
expect(searchAt).toBeGreaterThanOrEqual(0)
expect(chatAt).toBeGreaterThan(searchAt)
// Each id belongs to its own section.
expect(text.indexOf('aid-s')).toBeLessThan(chatAt)
expect(text.indexOf('aid-c')).toBeGreaterThan(chatAt)
})
it('shows only the configured default and says how many others exist', () => {
const entries = Array.from({ length: 5 }, (_x, i) => entry({ agent_id: `aid-${i}` }))
const text = catalog(entries, { defaultRetrieveAgentId: 'aid-3' }) ?? ''
expect(text).toContain('aid-3')
expect(text).not.toContain('aid-0')
expect(text).toContain('4 others exist')
})
it('falls back to the full list when the configured default is not in the cache', () => {
// A stale or mistyped default must not hide the services that do exist.
const text = catalog([entry({ agent_id: 'aid-1' })], { defaultRetrieveAgentId: 'aid-gone' }) ?? ''
expect(text).toContain('aid-1')
expect(text).not.toContain('default service')
})
it('caps the list at the limit, newest first, and states the shortfall', () => {
const entries = Array.from({ length: 14 }, (_x, i) => entry({
agent_id: `aid-${i}`,
// aid-13 newest, aid-0 oldest.
modify_time: `2026-08-${String(i + 1).padStart(2, '0')}T00:00:00`,
}))
const text = catalog(entries) ?? ''
expect(text).toContain('aid-13')
// The four oldest fall outside the window.
expect(text).not.toContain('aid-0 ')
expect(text).toContain(`showing ${CATALOG_ENTRY_LIMIT} most recently modified of 14`)
// Silent truncation would make the model treat the list as complete.
expect(text).toContain('--name <keyword>')
})
it('marks a capped fetch as a lower bound rather than an exact total', () => {
// The fetch itself stopped early, so even the count is unknown.
const entries = Array.from({ length: 12 }, (_x, i) => entry({ agent_id: `aid-${i}` }))
const text = catalog(entries, { truncated: true, total: 900 }) ?? ''
expect(text).toContain('more than')
})
it('states that agent_id is required and forbids guessing one', () => {
const text = catalog([entry({ agent_id: 'aid-1' })]) ?? ''
expect(text).toContain('required')
// Picking the closest-looking service returns unrelated evidence, which is
// worse for the user than an honest "no such knowledge base".
expect(text).toMatch(/say so plainly/i)
})
it('labels an unnamed service instead of rendering a bare dash', () => {
const text = catalog([entry({ agent_id: 'aid-1', agent_name: '' })]) ?? ''
expect(text).toContain('(unnamed)')
})
it('renders a description when one arrives, truncated to its budget', () => {
// Forward compatibility: the backend has not shipped this field yet.
const long = 'x'.repeat(250)
const text = catalog([entry({ agent_id: 'aid-1', description: long })]) ?? ''
expect(text).toContain('x'.repeat(199))
expect(text).not.toContain('x'.repeat(201))
expect(text).toContain('…')
})
})
@@ -0,0 +1,177 @@
import { describe, expect, it, vi } from 'vitest'
import { mkdtempSync } from 'node:fs'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import type { KbClient } from '../src/client.js'
import { ServiceCache, serviceCachePath, writeServiceCache } from '../src/service-cache.js'
import { installServiceContext } from '../src/service-context.js'
const HOST = 'cn-beijing.maas.aliyuncs.com'
const WS = 'llm-a'
const SOURCE_PLUGIN = 'tool-bailian-kb/services'
type Injected = { id: string; source: { plugin?: string }; content: { text?: string }[] }
type Decision = { kind: string; messages: Injected[] }
interface HarnessOptions {
entries?: { agent_id: string; agent_name: string; scene: 'search' | 'chat'; status: string }[]
cacheOverride?: ServiceCache
postJson?: () => Promise<unknown>
}
/** Drives the installed `agent/pre-step` listener against a fake session. */
function harness(options: HarnessOptions = {}) {
const home = mkdtempSync(join(tmpdir(), 'bailian-kb-ctx-'))
const entries = options.entries ?? [
{ agent_id: 'aid-1', agent_name: 'svc-one', scene: 'search' as const, status: 'deployed' },
]
writeServiceCache(serviceCachePath(WS, home), {
version: 1,
fetchedAt: Date.now(),
workspaceId: WS,
endpointHost: HOST,
entries,
total: entries.length,
truncated: false,
})
const postJson = vi.fn(options.postJson ?? (async () => ({ code: 'Success', data: { total_count: 0, rows: [] } })))
const cache = options.cacheOverride ?? new ServiceCache({
client: { postJson } as unknown as KbClient,
resolveWorkspaceId: async () => WS,
endpointHost: HOST,
warn: () => {},
home,
})
let listener!: (payload: unknown, next: () => Promise<unknown>) => Promise<unknown>
const warn = vi.fn()
const ctx = { on: (_event: string, handler: typeof listener) => { listener = handler } }
installServiceContext(ctx as never, {
cache,
resolveWorkspaceId: async () => WS,
resolveDefaultRetrieveAgentId: async () => undefined,
resolveDefaultChatAgentId: async () => undefined,
warn,
})
const events: { type: string; seq: number; data: unknown }[] = []
const agent = {
session: {
events,
surface: { nodes: [] as number[] },
header: { cwd: '/tmp' },
},
}
const run = async (proposed: Injected[], downstream: 'enter' | 'reject'): Promise<Decision> =>
await listener(
{ agent, signal: { aborted: false, throwIfAborted: () => {} } },
async () => downstream === 'reject' ? { kind: 'reject' } : { kind: 'enter', messages: proposed },
) as Decision
return {
cache,
warn,
postJson,
/** One step whose downstream decision enters with `proposed`. */
step: async (proposed: Injected[] = []) => await run(proposed, 'enter'),
/** One step whose downstream decision rejects. */
stepRejecting: async () => await run([], 'reject'),
/** Record an entered message into the durable log and make it visible. */
commit: (message: Injected) => {
events.push({ type: 'user/message', seq: events.length + 1, data: message })
agent.session.surface.nodes = events.map(event => event.seq)
},
/** Simulate compaction: the events remain, the surface shrinks. */
setSurface: (seqs: number[]) => { agent.session.surface.nodes = seqs },
/** This module's own messages within a decision. */
injected: (decision: Decision) => decision.messages.filter(m => m.source.plugin === SOURCE_PLUGIN),
}
}
describe('installServiceContext', () => {
it('injects the catalog as a sourced message on the first step', async () => {
const h = harness()
const decision = await h.step()
expect(decision.kind).toBe('enter')
const own = h.injected(decision)
expect(own).toHaveLength(1)
expect(own[0]?.content[0]?.text).toContain('aid-1')
})
it('does not re-inject on later steps of the same turn', async () => {
// pre-step fires once per model request, so a turn with five tool calls
// fires it six times. Re-injecting each time would insert six copies and
// void the KV cache from the first insertion onward — this is a correctness
// guard, not an optimization.
const h = harness()
const first = await h.step()
const own = h.injected(first)
expect(own).toHaveLength(1)
h.commit(own[0]!)
for (let toolRoundTrip = 0; toolRoundTrip < 5; toolRoundTrip += 1) {
expect(h.injected(await h.step())).toHaveLength(0)
}
})
it('re-injects once compaction drops the catalog from the visible surface', async () => {
const h = harness()
const first = await h.step()
h.commit(h.injected(first)[0]!)
expect(h.injected(await h.step())).toHaveLength(0)
// The event stays in the durable log but leaves the surface. Comparing only
// "did we ever publish this" would suppress every future injection and the
// model would finish the session with no service list at all.
h.setSurface([])
expect(h.injected(await h.step())).toHaveLength(1)
})
it('replaces its own pending message instead of adding a second one', async () => {
const h = harness()
const pending = h.injected(await h.step())[0]!
const second = await h.step([pending])
expect(h.injected(second)).toHaveLength(1)
})
it('injects nothing when the cache holds no services', async () => {
const h = harness({ entries: [] })
expect(h.injected(await h.step())).toHaveLength(0)
})
it('never awaits the network refresh', async () => {
// A slow list request must not delay the user's request.
let settle: ((value: unknown) => void) | undefined
const h = harness({ postJson: async () => await new Promise(resolve => { settle = resolve }) })
// Force staleness so this step schedules a refresh.
h.cache.invalidate()
const raced = await Promise.race([
h.step().then(() => 'stepped'),
new Promise(resolve => { setTimeout(() => resolve('timed out'), 200) }),
])
expect(raced).toBe('stepped')
settle?.({ data: { rows: [] } })
})
it('degrades to no injection when anything throws, never failing the step', async () => {
// A throwing pre-step listener fails the proposed step, i.e. stalls the
// user's turn. A missing catalog is far cheaper than that.
const exploding = {
isStale: () => { throw new Error('corrupt cache') },
peek: () => undefined,
refresh: async () => {},
invalidate: () => {},
} as unknown as ServiceCache
const h = harness({ cacheOverride: exploding })
const decision = await h.step()
expect(decision.kind).toBe('enter')
expect(h.injected(decision)).toHaveLength(0)
expect(h.warn).toHaveBeenCalled()
})
it('passes a rejected downstream decision straight through', async () => {
const h = harness()
expect((await h.stepRejecting()).kind).toBe('reject')
})
})
@@ -0,0 +1,119 @@
import { describe, expect, it, vi } from 'vitest'
import type { ServiceListResponse } from '../src/api-types.js'
import type { KbClient } from '../src/client.js'
import { listServices } from '../src/services.js'
/** A page of `count` rows, all deployed unless overridden. */
function page(count: number, total: number, overrides: Record<string, unknown> = {}): ServiceListResponse {
return {
code: 'Success',
data: {
total_count: total,
rows: Array.from({ length: count }, (_row, index) => ({
agent_id: `aid-${index}`,
agent_name: `service-${index}`,
agent_status: 'deployed',
modify_time: '2026-08-20T10:00:00',
...overrides,
})),
},
}
}
/** A client whose postJson is driven by a queue of per-call responses or errors. */
function clientReturning(...outcomes: (ServiceListResponse | Error)[]): {
client: KbClient
postJson: ReturnType<typeof vi.fn>
} {
const postJson = vi.fn(async () => {
const next = outcomes.shift()
if (next === undefined) throw new Error('unexpected extra request')
if (next instanceof Error) throw next
return next
})
return { client: { postJson } as unknown as KbClient, postJson }
}
describe('listServices', () => {
it('queries both scenes for deployed services only and tags each entry with its scene', async () => {
const { client, postJson } = clientReturning(page(1, 1), page(1, 1))
const result = await listServices(client)
expect(postJson.mock.calls[0]?.[1]).toMatchObject({
agent_scene: 'search',
// Verified server-side: this filter is honored and excludes drafts.
agent_status: 'deployed',
page_number: 1,
page_size: 100,
})
expect(postJson.mock.calls[1]?.[1]).toMatchObject({ agent_scene: 'chat' })
expect(result.entries.map(e => e.scene)).toEqual(['search', 'chat'])
expect(result.total).toBe(2)
expect(result.truncated).toBe(false)
expect(result.failedScenes).toEqual([])
})
it('stops at a short page without asking for another', async () => {
// 3 rows on a 100-row page is the last page; a second request would be waste.
const { client, postJson } = clientReturning(page(3, 3), page(0, 0))
const result = await listServices(client)
expect(postJson).toHaveBeenCalledTimes(2) // one per scene, not one per page
expect(result.entries).toHaveLength(3)
})
it('caps at two pages per scene and reports the shortfall as truncated', async () => {
// A workspace claiming 500 rows: fetch 200, flag the rest as unfetched.
const { client, postJson } = clientReturning(
page(100, 500),
page(100, 500),
page(0, 0),
)
const result = await listServices(client)
// 2 pages for search + 1 short page for chat: the cap holds.
expect(postJson).toHaveBeenCalledTimes(3)
expect(result.entries).toHaveLength(200)
expect(result.truncated).toBe(true)
})
it('keeps one scene when the other fails instead of losing the whole list', async () => {
const { client } = clientReturning(page(2, 2), new Error('chat scene exploded'))
const result = await listServices(client)
expect(result.entries).toHaveLength(2)
expect(result.entries.every(e => e.scene === 'search')).toBe(true)
expect(result.failedScenes).toEqual(['chat'])
})
it('reports both scenes as failed without throwing', async () => {
const { client } = clientReturning(new Error('down'), new Error('down'))
const result = await listServices(client)
expect(result.entries).toEqual([])
expect(result.failedScenes).toEqual(['search', 'chat'])
})
it('drops rows without an agent_id and never reads pipeline_list', async () => {
// pipeline_list is unreliable in production (missing names, sometimes empty),
// so entries must not carry any knowledge-base label derived from it.
const { client } = clientReturning(
{
code: 'Success',
data: {
total_count: 2,
rows: [
{ agent_name: 'no id', agent_status: 'deployed' },
{ agent_id: 'aid-1', agent_name: 'ok', agent_status: 'deployed', pipeline_list: [{ pipeline_id: 'p1' }] },
],
},
},
page(0, 0),
)
const result = await listServices(client)
expect(result.entries).toHaveLength(1)
expect(result.entries[0]).toEqual({
agent_id: 'aid-1',
agent_name: 'ok',
scene: 'search',
status: 'deployed',
})
expect(JSON.stringify(result.entries)).not.toContain('p1')
})
})
@@ -0,0 +1,93 @@
import { describe, expect, it } from 'vitest'
import { readFileSync } from 'node:fs'
import { fileURLToPath } from 'node:url'
import { parseSkillFile } from '../src/skill.js'
const SKILL_PATH = fileURLToPath(new URL('../skills/bailian-kb/SKILL.md', import.meta.url))
describe('parseSkillFile', () => {
it('takes name and description from frontmatter and strips the block from the body', () => {
const parsed = parseSkillFile([
'---',
'name: demo-skill',
'description: >-',
' First line of the folded description,',
' continued on the next source line.',
'---',
'',
'# Heading',
'',
'Body text.',
'',
].join('\n'))
expect(parsed?.name).toBe('demo-skill')
// A folded scalar joins its lines with spaces; the value must arrive unfolded.
expect(parsed?.description).toBe('First line of the folded description, continued on the next source line.')
// The registry performs no parsing of its own, so the YAML must already be gone.
expect(parsed?.content).toBe('\n# Heading\n\nBody text.\n')
expect(parsed?.content).not.toContain('---')
expect(parsed?.content).not.toContain('description:')
})
it('carries optional whenToUse and metadata, omitting them when absent or blank', () => {
const withExtras = parseSkillFile([
'---',
'name: demo-skill',
'description: Routing text.',
'whenToUse: Extra guidance.',
'metadata:',
' owner: platform',
'---',
'Body.',
].join('\n'))
expect(withExtras?.whenToUse).toBe('Extra guidance.')
expect(withExtras?.metadata).toEqual({ owner: 'platform' })
const bare = parseSkillFile('---\nname: demo-skill\ndescription: Routing text.\n---\nBody.')
expect(bare).not.toHaveProperty('whenToUse')
expect(bare).not.toHaveProperty('metadata')
})
it('rejects a file without usable frontmatter instead of throwing', () => {
// No fence at all: a plain markdown file.
expect(parseSkillFile('# Just markdown\n')).toBeUndefined()
// Opening fence never closes.
expect(parseSkillFile('---\nname: demo-skill\n')).toBeUndefined()
// Unparsable YAML inside the fence.
expect(parseSkillFile('---\nname: [unclosed\n---\nBody.')).toBeUndefined()
// Scalar and sequence roots are not frontmatter records.
expect(parseSkillFile('---\njust a string\n---\nBody.')).toBeUndefined()
expect(parseSkillFile('---\n- one\n- two\n---\nBody.')).toBeUndefined()
// name or description missing, blank, or the wrong type.
expect(parseSkillFile('---\ndescription: Routing text.\n---\nBody.')).toBeUndefined()
expect(parseSkillFile('---\nname: demo-skill\n---\nBody.')).toBeUndefined()
expect(parseSkillFile('---\nname: demo-skill\ndescription: " "\n---\nBody.')).toBeUndefined()
expect(parseSkillFile('---\nname: 42\ndescription: Routing text.\n---\nBody.')).toBeUndefined()
})
})
describe('the packaged bailian-kb SKILL.md', () => {
const parsed = parseSkillFile(readFileSync(SKILL_PATH, 'utf8'))
it('parses, and is the single source of the registered name and description', () => {
expect(parsed?.name).toBe('bailian-kb')
expect(parsed?.description).toContain('bl')
// The frontmatter must state where credentials come from: that sentence used
// to live in skill.ts and would otherwise be lost with the duplicate.
expect(parsed?.description).toContain('DASHSCOPE_API_KEY')
})
it('keeps its description within the catalog truncation budget', () => {
// tool-skill publishes catalog entries through catalogDescription(), which
// whitespace-normalizes and hard-truncates at 500 characters.
const normalized = (parsed?.description ?? '').replaceAll(/\s+/g, ' ').trim()
expect(normalized.length).toBeLessThanOrEqual(500)
})
it('registers a body with no frontmatter residue', () => {
expect(parsed?.content.startsWith('---')).toBe(false)
expect(parsed?.content).not.toContain('description: >-')
// The real body still begins with the document heading.
expect(parsed?.content.trimStart().startsWith('#')).toBe(true)
})
})
+53 -3
View File
@@ -4,9 +4,9 @@ import { createKbTools } from '../src/tools.js'
const EXEC = {} as never
function toolsWith(postJson: unknown, postSse?: unknown, resolveDefaultRetrieveAgentId?: () => Promise<string | undefined>, resolveDefaultChatAgentId?: () => Promise<string | undefined>) {
function toolsWith(postJson: unknown, postSse?: unknown, resolveDefaultRetrieveAgentId?: () => Promise<string | undefined>, resolveDefaultChatAgentId?: () => Promise<string | undefined>, describeServicesAfterRefresh?: (scene: 'search' | 'chat') => Promise<string | undefined>) {
const client = { postJson, postSse, agentVersion: undefined } as unknown as KbClient
const list = createKbTools({ client, ...(resolveDefaultRetrieveAgentId ? { resolveDefaultRetrieveAgentId } : {}), ...(resolveDefaultChatAgentId ? { resolveDefaultChatAgentId } : {}), chatTimeoutMs: 1000 })
const list = createKbTools({ client, ...(resolveDefaultRetrieveAgentId ? { resolveDefaultRetrieveAgentId } : {}), ...(resolveDefaultChatAgentId ? { resolveDefaultChatAgentId } : {}), ...(describeServicesAfterRefresh ? { describeServicesAfterRefresh } : {}), chatTimeoutMs: 1000 })
const byName = Object.fromEntries(list.map(t => [t.name, t]))
return { byName, list }
}
@@ -26,6 +26,17 @@ describe('createKbTools', () => {
expect(list.map(t => t.name).sort()).toEqual(['kb_chat', 'kb_search'])
})
it('keeps both descriptions free of service ids so the schemas stay prefix-stable', () => {
// The live service list rides an `agent/pre-step` context message precisely
// because re-registering a tool to refresh its description would void the
// prompt prefix cache. Any id leaking in here means that decision regressed.
const { byName } = toolsWith(vi.fn())
for (const tool of [byName.kb_search!, byName.kb_chat!]) {
const text = `${tool.description} ${JSON.stringify(tool.parameters)}`
expect(text).not.toMatch(/aid-[0-9a-f]/)
}
})
it('kb_search truncates nodes client-side to top_k and never sends top_k to the server', async () => {
const postJson = vi.fn(async (_path: string, _body: unknown) => searchResponse)
const { byName } = toolsWith(postJson)
@@ -73,7 +84,7 @@ describe('createKbTools', () => {
expect((postJson.mock.calls[0]![1] as Record<string, unknown>).agent_id).toBe('aid-explicit')
})
it('a 4xx failure passes the original error through unchanged', async () => {
it('a 4xx failure passes the original error through unchanged when no refresh hook is wired', async () => {
const postJson = vi.fn(async (_path: string) => {
throw new KbApiError('agent not found', 400)
})
@@ -82,6 +93,45 @@ describe('createKbTools', () => {
expect((err as Error).message).toBe('agent not found')
})
it('a 4xx failure appends the refreshed service list for the calling scene', async () => {
// agent_id is the only caller-supplied identifier on these endpoints, so a
// rejected request most often means the cached id is gone. The recovery the
// model needs is the current list, and an error message carries it without
// disturbing the request prefix.
const postJson = vi.fn(async (_path: string) => {
throw new KbApiError('agent not found', 400)
})
const describe = vi.fn(async (scene: 'search' | 'chat') => `services for ${scene}: aid-new`)
const { byName } = toolsWith(postJson, undefined, undefined, undefined, describe)
const err = await byName.kb_search!.execute({ query: 'q', agent_id: 'stale' }, EXEC).catch((e: unknown) => e)
expect(describe).toHaveBeenCalledWith('search')
expect((err as Error).message).toContain('agent not found')
expect((err as Error).message).toContain('services for search: aid-new')
})
it('leaves a 5xx failure and a failing refresh alone', async () => {
// A server-side fault is not an id problem, and a refresh that itself fails
// must not replace the real error with its own.
const serverError = vi.fn(async (_path: string) => {
throw new KbApiError('upstream exploded', 502)
})
const describe = vi.fn(async () => 'never used')
const { byName } = toolsWith(serverError, undefined, undefined, undefined, describe)
const err = await byName.kb_search!.execute({ query: 'q', agent_id: 'aid-1' }, EXEC).catch((e: unknown) => e)
expect(describe).not.toHaveBeenCalled()
expect((err as Error).message).toBe('upstream exploded')
const badRequest = vi.fn(async (_path: string) => {
throw new KbApiError('agent not found', 400)
})
const failing = vi.fn(async () => {
throw new Error('refresh also down')
})
const { byName: byName2 } = toolsWith(badRequest, undefined, undefined, undefined, failing)
const err2 = await byName2.kb_search!.execute({ query: 'q', agent_id: 'bad' }, EXEC).catch((e: unknown) => e)
expect((err2 as Error).message).toBe('agent not found')
})
it('kb_chat buffers the SSE stream into one answer', async () => {
const sse = 'data: {"output":{"choices":[{"message":{"content":"hi"},"finish_reason":"stop"}]},"request_id":"r2"}\n\ndata: [DONE]\n\n'
const postSse = vi.fn(async () => new Response(sse, { status: 200 }))
+31 -10
View File
@@ -13,13 +13,20 @@ importers:
version: 5.9.3
vitest:
specifier: ^3.0.0
version: 3.2.7(@types/debug@4.1.13)(@types/node@22.20.1)(lightningcss@1.33.0)
version: 3.2.7(@types/debug@4.1.13)(@types/node@22.20.1)(lightningcss@1.33.0)(yaml@2.9.0)
packages/tool-bailian-kb:
dependencies:
yaml:
specifier: ^2.4.2
version: 2.9.0
devDependencies:
'@deepseek-ai/cordis':
specifier: ^4.0.1
version: 4.0.1(@deepseek-ai/cordis-plugin-include@1.0.6)(@deepseek-ai/cordis-plugin-loader@1.0.2)
'@deepseek-ai/dsh-agent':
specifier: ^0.1.0-rc.6
version: 0.1.0-rc.6(b0514d7a320728b6d8f5a31eb3960e10)
'@deepseek-ai/dsh-api-remotes':
specifier: ^0.1.0-rc.6
version: 0.1.0-rc.6(963f21f11cddd09def26a2473cddd8c5)
@@ -41,6 +48,12 @@ importers:
'@deepseek-ai/dsh-credentials':
specifier: ^0.1.0-rc.6
version: 0.1.0-rc.6(@deepseek-ai/cordis@4.0.1)(@deepseek-ai/dsh-brand@0.1.0-rc.6(@deepseek-ai/cordis@4.0.1)(@deepseek-ai/dsh-invariants@0.1.0-rc.6(@deepseek-ai/cordis@4.0.1)))(@deepseek-ai/dsh-invariants@0.1.0-rc.6(@deepseek-ai/cordis@4.0.1))
'@deepseek-ai/dsh-llm':
specifier: ^0.1.0-rc.6
version: 0.1.0-rc.6(@deepseek-ai/cordis@4.0.1)(@deepseek-ai/dsh-attachment@0.1.0-rc.6(@deepseek-ai/cordis@4.0.1)(@deepseek-ai/dsh-brand@0.1.0-rc.6(@deepseek-ai/cordis@4.0.1)(@deepseek-ai/dsh-invariants@0.1.0-rc.6(@deepseek-ai/cordis@4.0.1)))(@deepseek-ai/dsh-invariants@0.1.0-rc.6(@deepseek-ai/cordis@4.0.1)))(@deepseek-ai/dsh-brand@0.1.0-rc.6(@deepseek-ai/cordis@4.0.1)(@deepseek-ai/dsh-invariants@0.1.0-rc.6(@deepseek-ai/cordis@4.0.1)))(@deepseek-ai/dsh-invariants@0.1.0-rc.6(@deepseek-ai/cordis@4.0.1))(@deepseek-ai/dsh-timeout@0.1.0-rc.6(@deepseek-ai/cordis@4.0.1)(@deepseek-ai/dsh-invariants@0.1.0-rc.6(@deepseek-ai/cordis@4.0.1)))
'@deepseek-ai/dsh-session':
specifier: ^0.1.0-rc.6
version: 0.1.0-rc.6(6fd26f59436a18b115f326d6060415e6)
'@deepseek-ai/dsh-settings':
specifier: ^0.1.0-rc.6
version: 0.1.0-rc.6(@deepseek-ai/cordis@4.0.1)(@deepseek-ai/dsh-brand@0.1.0-rc.6(@deepseek-ai/cordis@4.0.1)(@deepseek-ai/dsh-invariants@0.1.0-rc.6(@deepseek-ai/cordis@4.0.1)))(@deepseek-ai/dsh-invariants@0.1.0-rc.6(@deepseek-ai/cordis@4.0.1))(@deepseek-ai/schemastery@3.18.1)
@@ -1935,6 +1948,11 @@ packages:
utf-8-validate:
optional: true
yaml@2.9.0:
resolution: {integrity: sha512-2AvhNX3mb8zd6Zy7INTtSpl1F15HW6Wnqj0srWlkKLcpYl/gMIMJiyuGq2KeI2YFxUPjdlB+3Lc10seMLtL4cA==}
engines: {node: '>= 14.6'}
hasBin: true
yuku-ast@0.8.7:
resolution: {integrity: sha512-h6+4bDfyootiMB9vckk5uKo5r5j0GHrkr17FQTDNfEsFT3DWlN9uu1HJwQwc64pgmLCI945fWM3lbTIqxjT3GQ==}
@@ -2867,13 +2885,13 @@ snapshots:
chai: 5.3.3
tinyrainbow: 2.0.0
'@vitest/mocker@3.2.7(vite@7.3.6(@types/node@22.20.1)(lightningcss@1.33.0))':
'@vitest/mocker@3.2.7(vite@7.3.6(@types/node@22.20.1)(lightningcss@1.33.0)(yaml@2.9.0))':
dependencies:
'@vitest/spy': 3.2.7
estree-walker: 3.0.3
magic-string: 0.30.21
optionalDependencies:
vite: 7.3.6(@types/node@22.20.1)(lightningcss@1.33.0)
vite: 7.3.6(@types/node@22.20.1)(lightningcss@1.33.0)(yaml@2.9.0)
'@vitest/pretty-format@3.2.7':
dependencies:
@@ -3770,13 +3788,13 @@ snapshots:
'@types/unist': 3.0.3
vfile-message: 4.0.3
vite-node@3.2.4(@types/node@22.20.1)(lightningcss@1.33.0):
vite-node@3.2.4(@types/node@22.20.1)(lightningcss@1.33.0)(yaml@2.9.0):
dependencies:
cac: 6.7.14
debug: 4.4.3
es-module-lexer: 1.7.0
pathe: 2.0.3
vite: 7.3.6(@types/node@22.20.1)(lightningcss@1.33.0)
vite: 7.3.6(@types/node@22.20.1)(lightningcss@1.33.0)(yaml@2.9.0)
transitivePeerDependencies:
- '@types/node'
- jiti
@@ -3791,7 +3809,7 @@ snapshots:
- tsx
- yaml
vite@7.3.6(@types/node@22.20.1)(lightningcss@1.33.0):
vite@7.3.6(@types/node@22.20.1)(lightningcss@1.33.0)(yaml@2.9.0):
dependencies:
esbuild: 0.28.2
fdir: 6.5.0(picomatch@4.0.5)
@@ -3803,12 +3821,13 @@ snapshots:
'@types/node': 22.20.1
fsevents: 2.3.3
lightningcss: 1.33.0
yaml: 2.9.0
vitest@3.2.7(@types/debug@4.1.13)(@types/node@22.20.1)(lightningcss@1.33.0):
vitest@3.2.7(@types/debug@4.1.13)(@types/node@22.20.1)(lightningcss@1.33.0)(yaml@2.9.0):
dependencies:
'@types/chai': 5.2.3
'@vitest/expect': 3.2.7
'@vitest/mocker': 3.2.7(vite@7.3.6(@types/node@22.20.1)(lightningcss@1.33.0))
'@vitest/mocker': 3.2.7(vite@7.3.6(@types/node@22.20.1)(lightningcss@1.33.0)(yaml@2.9.0))
'@vitest/pretty-format': 3.2.7
'@vitest/runner': 3.2.7
'@vitest/snapshot': 3.2.7
@@ -3826,8 +3845,8 @@ snapshots:
tinyglobby: 0.2.17
tinypool: 1.1.1
tinyrainbow: 2.0.0
vite: 7.3.6(@types/node@22.20.1)(lightningcss@1.33.0)
vite-node: 3.2.4(@types/node@22.20.1)(lightningcss@1.33.0)
vite: 7.3.6(@types/node@22.20.1)(lightningcss@1.33.0)(yaml@2.9.0)
vite-node: 3.2.4(@types/node@22.20.1)(lightningcss@1.33.0)(yaml@2.9.0)
why-is-node-running: 2.3.0
optionalDependencies:
'@types/debug': 4.1.13
@@ -3853,6 +3872,8 @@ snapshots:
ws@8.21.3: {}
yaml@2.9.0: {}
yuku-ast@0.8.7:
dependencies:
'@yuku-toolchain/types': 0.8.7