Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
Next Next commit
fix(runtime-broker): opt in to v2 execution reconciliation
  • Loading branch information
wenshao committed Sep 29, 2026
commit 7444dc9726e1bc8ce8fdf8d15093b2506e2f2abd
6 changes: 3 additions & 3 deletions docs/design/2026-09-28-managed-mcp-runtime.md
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,7 @@ Set `QWEN_MANAGED_MCP_CONFIG` on the Runtime process to an absolute manifest pat
}
```

For `streamable-http` or `sse`, use `url` and optional `headers` instead of `command`, `args` and `env`. Definitions are loaded at Runtime startup; a changed recipe requires a new `serverRevision` and matching digest. Keep both revisions available when replacing a live binding. Optional `timeoutMs` is an integer from 1 to 600000 and defaults to 600000 for invocation responses. Connect/discovery retain their 25-second bound. Hosted tool observation lasts 630 seconds, queries the original execution through temporary UNKNOWN responses, and Java persists a conclusive late Runtime result for both tool protocol versions. Polling reuses the acquired owner; failed status queries re-establish only that original owner.
For `streamable-http` or `sse`, use `url` and optional `headers` instead of `command`, `args` and `env`. Definitions are loaded at Runtime startup; a changed recipe requires a new `serverRevision` and matching digest. Keep both revisions available when replacing a live binding. Optional `timeoutMs` is an integer from 1 to 600000 and defaults to 600000 for invocation responses. Connect/discovery retain their 25-second bound. Hosted tool observation lasts 630 seconds, queries the original execution through temporary UNKNOWN responses, and MCP polling explicitly opts into original-execution reconciliation with `GET /executions/:id?reconcile=true`; Java persists only a conclusive late Runtime result. V2 status reads without this opt-in remain passive, while v3 retains its existing automatic reconciliation. The optional query accepts `true` or `false` and does not change execution ownership or permit dispatch. Polling reuses the acquired owner; failed status queries re-establish only that original owner.

Create or load a private Hosted Session with its usual `managedSessionStore` fields plus `toolProfile: "hosted-workspace-mcp/1"` and `mcpServers: [{serverId, serverRevision, definitionDigest}]`. The Session retains these initial admission pins; replacements are committed separately. The first prompt or resource/prompt operation initializes the bindings before admitting the work. Hosted requests use the existing Harness protocol/boot and client identity headers.

Expand All @@ -101,12 +101,12 @@ Required checks include Session/workspace isolation, credential omission, succes

H1 deliberately permits only one attached MCP owner per tenant/storage lease. This excludes other Sessions in the same Workspace and also other Workspaces sharing that storage, even between turns. Detach releases the lease. Releasing it while a stdio server still has Workspace access would allow concurrent writers; separating connection lifetime from storage ownership is follow-up work before production enablement.

Physical connection loss with an unresolved request remains recovery-blocked. A tool that settles after the 630-second Hosted observation window, or after Harness restart during its turn, still needs checkpoint recovery that this private slice does not implement. Raw resource/prompt operations can accept a late result through their original status route or close. Runtime receipt history and closed connection tombstones remain generation-local and grow for the process lifetime; pruning with durable acknowledgement, as well as eliminating the pending-release waiter on a permanently lost request, remains follow-up work. Do not treat the concurrency quotas as a bound on history memory.
Physical connection loss with an unresolved request remains recovery-blocked. A lost connection or a server that never replies can consume the full 630-second Hosted observation window; a shorter Runtime timeout or cancellation does not shorten that window. Every pinned server must refresh successfully before the model request, so an unavailable server also blocks text-only turns and calls to healthy servers until it recovers. A tool that settles after the 630-second Hosted observation window, or after Harness restart during its turn, still needs checkpoint recovery that this private slice does not implement. Raw resource/prompt operations can accept a late result through their original status route or close. Runtime receipt history and closed connection tombstones remain generation-local and grow for the process lifetime; pruning with durable acknowledgement, as well as eliminating the pending-release waiter on a permanently lost request, remains follow-up work. Do not treat the concurrency quotas as a bound on history memory.

Model tool names intentionally include catalog and connection identity so old advertised calls cannot silently use a newer binding. A reloaded transcript can contain historical names. The stdio HOME/USERPROFILE directory is the Workspace; deployment definitions should explicitly set a separate HOME if a server writes caches there.

The record domains are globally recognized by the Session store, while execution remains gated by the explicit private profile and scoped Runtime definitions. Shared Broker acquire is intentionally idempotent for its existing owner and returns workspace generation plus Runtime binding/generation; ordinary file/Shell prepare calls keep the actual prompt and call identity.

The unreleased MCP migration is V21 to avoid the V19/V20 publication migrations on #12894. Deploy both migrations in increasing order; whichever branch lands later must recheck against main. The generic controls in #12868 require a semantic merge that preserves MCP original-owner recovery after revocation and drain-before-storage-release ordering.
The unreleased MCP migration is V21 to avoid the V19/V20 publication migrations on #12894. Deploy both migrations in increasing order; whichever branch lands later must recheck against main. If MCP V21 has already run before the publication migrations arrive, those pending migrations must be renumbered above the deployed version; do not enable out-of-order migration to bypass that check. The generic controls in #12868 require a semantic merge that preserves MCP original-owner recovery after revocation and drain-before-storage-release ordering.

Public configuration management and production AgentBundle capability publication remain separate deployments. Remote systems without query or idempotency support cannot recover an unknown effect automatically. The Runtime's generation-local receipts are not durable across physical Runtime loss; committed Session intents preserve the blocked outcome in that case. Quotas are per Runtime instance, not an aggregate across separate Runtime processes. The private profile uses the existing inline Session Store: each discovery list is limited to 16 KiB, with retained prefixes marked partial and no usable entries marked failed; raw operation responses are limited to 60 KiB and oversized responses settle with an output-limit error. SDK reverse clients, production profile advertisement, cross-process aggregate budgets and object-storage results are not enabled by this slice.
6 changes: 3 additions & 3 deletions docs/design/2026-09-28-managed-mcp-runtime.zh-CN.md
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,7 @@ Runtime 只允许目标 workspace 已配置的定义,以及 Session 已安装
}
```

`streamable-http` 或 `sse` 使用 `url` 和可选 `headers`,不能同时提供 `command`、`args`、`env`。定义在 Runtime 启动时装载;修改连接配方必须提供新的 `serverRevision` 和对应 digest。替换活动 binding 时应同时保留两个修订。可选 `timeoutMs` 为 1 至 600000 的整数,调用响应默认等待 600000 毫秒;连接和发现仍为 25 秒上限。Hosted 工具观察窗口为 630 秒,在临时 UNKNOWN 后继续查询原 execution;Java 对两种工具协议均持久接受明确的迟到 Runtime 结果。轮询复用已取得的 owner,状态查询失败时仅重建原 owner 路由。
`streamable-http` 或 `sse` 使用 `url` 和可选 `headers`,不能同时提供 `command`、`args`、`env`。定义在 Runtime 启动时装载;修改连接配方必须提供新的 `serverRevision` 和对应 digest。替换活动 binding 时应同时保留两个修订。可选 `timeoutMs` 为 1 至 600000 的整数,调用响应默认等待 600000 毫秒;连接和发现仍为 25 秒上限。Hosted 工具观察窗口为 630 秒,在临时 UNKNOWN 后继续查询原 execution;MCP 轮询通过 `GET /executions/:id?reconcile=true` 显式请求原 execution 对账,Java 仅持久接受明确的迟到 Runtime 结果。未显式请求对账的 v2 状态查询保持被动读取,v3 保留既有自动对账。可选查询参数接受 `true` 或 `false`,不改变 execution 归属,也不允许派发。轮询复用已取得的 owner,状态查询失败时仅重建原 owner 路由。

创建或加载私有 Hosted Session 时,在原有 `managedSessionStore` 字段之外提供 `toolProfile: "hosted-workspace-mcp/1"` 和 `mcpServers: [{serverId, serverRevision, definitionDigest}]`。Session 保留初始准入 pin,后续替换单独提交。首个 prompt 或资源/提示操作会先初始化 binding,再受理工作。Hosted 请求沿用现有 Harness protocol/boot 与 client 身份 headers。

Expand All @@ -101,12 +101,12 @@ Runtime 只允许目标 workspace 已配置的定义,以及 Session 已安装

H1 明确保留每个 tenant/storage lease 同时仅一个 attached MCP owner 的限制。同一 Workspace 的其他 Session,以及共享该 storage 的其他 Workspace,即使在两次 turn 之间也不能执行;detach 后释放租约。stdio server 仍能访问 Workspace 时提前释放租约会允许并发写入;连接寿命与存储所有权解耦属于生产启用前的后续工作。

存在未决请求时物理连接丢失仍保持 recovery-blocked。工具在 630 秒 Hosted 观察窗口结束后才结算,或 Harness 在 turn 中途重启,仍需要本私有阶段尚未实现的 checkpoint 恢复。原始资源/提示操作可通过原 ID 的状态接口或 close 接受迟到结果。Runtime 回执历史及已关闭连接的 tombstone 仍保留至进程结束,随历史增长;基于持久 ACK 的回收,以及永久丢失请求的 release 等待者清理,留作后续工作。并发配额不代表历史内存有上限。
存在未决请求时物理连接丢失仍保持 recovery-blocked。连接丢失或 server 永不回复可能耗尽整个 630 秒 Hosted 观察窗口;更短的 Runtime 超时或取消不会缩短该窗口。模型请求前必须成功刷新所有固定 server,因此一个 server 不可用也会阻塞纯文本 turn 和健康 server 的调用,直到它恢复。工具在 630 秒 Hosted 观察窗口结束后才结算,或 Harness 在 turn 中途重启,仍需要本私有阶段尚未实现的 checkpoint 恢复。原始资源/提示操作可通过原 ID 的状态接口或 close 接受迟到结果。Runtime 回执历史及已关闭连接的 tombstone 仍保留至进程结束,随历史增长;基于持久 ACK 的回收,以及永久丢失请求的 release 等待者清理,留作后续工作。并发配额不代表历史内存有上限。

模型工具名有意包含目录和连接身份,防止旧广告调用静默使用新 binding;重新加载的历史可能保留旧名字。stdio HOME/USERPROFILE 为 Workspace;会写 HOME 缓存的 server 应通过部署定义显式设置独立 HOME。

Session 存储全局识别这两个记录 domain;实际执行仍由显式私有 profile 和限定范围的 Runtime 定义控制。共享 Broker acquire 有意对原 owner 幂等,并返回 workspace generation 与 Runtime binding/generation;普通文件/Shell prepare 保留实际 prompt 和 call 身份。

尚未发布的 MCP migration 使用 V21,避免与 #12894 的 V19/V20 publication migration 重号。部署应按递增顺序执行,后合并分支必须再次对照 main 检查。#12868 的通用 control 需要语义合并,保留撤权后的 MCP 原 owner 恢复,以及先 drain 再释放存储租约的顺序。
尚未发布的 MCP migration 使用 V21,避免与 #12894 的 V19/V20 publication migration 重号。部署应按递增顺序执行,后合并分支必须再次对照 main 检查。如果 MCP V21 已先执行,后到达且尚未应用的 publication migration 必须重新编号到已部署版本之后,不能靠启用 out-of-order migration 绕过检查。#12868 的通用 control 需要语义合并,保留撤权后的 MCP 原 owner 恢复,以及先 drain 再释放存储租约的顺序。

公开配置管理和生产 AgentBundle 能力发布另行部署。没有查询或幂等支持的远端系统不能自动恢复未知副作用。Runtime 的代内回执不能跨物理 Runtime 丢失持久保留;此时已提交的 Session intent 保持阻塞结果。配额按 Runtime 实例计算,不跨独立 Runtime 进程汇总。私有配置沿用现有 inline Session Store:每类发现列表限制为 16 KiB,保留部分条目时标记 partial,无法保留有效条目时标记 failed;原始操作响应限制为 60 KiB,超限以 output-limit 错误结算。本阶段不启用 SDK 反向客户端、生产 profile 公告、跨进程总预算或对象存储结果。
18 changes: 14 additions & 4 deletions packages/cli/src/serve/hosted-workspace-broker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,9 +34,10 @@ async function fixture(
const chunks: Buffer[] = [];
for await (const chunk of req) chunks.push(Buffer.from(chunk));
const body = Buffer.concat(chunks).toString();
const url = new URL(req.url!, 'http://fixture');
const response = handler(
new URL(req.url!, 'http://fixture').pathname,
body ? JSON.parse(body) : {},
url.pathname,
body ? JSON.parse(body) : Object.fromEntries(url.searchParams),
);
if (response.drop) {
res.destroy();
Expand Down Expand Up @@ -317,8 +318,9 @@ it.each(['runtime_idempotency_conflict', 'runtime_execution_conflict'])(

it('queries the original identity when start reports an unknown execution', async () => {
const paths: string[] = [];
const broker = await fixture((path) => {
const broker = await fixture((path, fields) => {
paths.push(path);
expect(fields).not.toHaveProperty('reconcile');
return { code: 409, body: { code: 'runtime_broker_execution_unknown' } };
});
await expect(
Expand All @@ -332,9 +334,11 @@ it('queries the original identity when start reports an unknown execution', asyn

it('observes a late original result after unknown cancellation without starting again', async () => {
const paths: string[] = [];
const queries: Array<Record<string, unknown>> = [];
let observations = 0;
const broker = await fixture((path) => {
const broker = await fixture((path, fields) => {
paths.push(path);
if (!path.endsWith(':cancel')) queries.push(fields);
if (path.endsWith(':cancel') || observations++ === 0)
return { code: 409, body: { code: 'runtime_broker_execution_unknown' } };
return {
Expand All @@ -356,6 +360,12 @@ it('observes a late original result after unknown cancellation without starting
expect(paths.filter((path) => path.endsWith(':cancel'))).toHaveLength(1);
expect(paths.filter((path) => path.endsWith(':start'))).toHaveLength(0);
expect(paths.filter((path) => path.endsWith('/execution'))).toHaveLength(2);
for (const query of queries)
expect(query).toMatchObject({
reconcile: 'true',
harnessSessionId: 'session',
runtimeSessionId: 'turn',
});
});

it('queries the original identity after an uncertain start failure', async () => {
Expand Down
4 changes: 3 additions & 1 deletion packages/cli/src/serve/hosted-workspace-broker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -194,7 +194,9 @@ export class HostedWorkspaceBroker {
cancellationSent = true;
response = await this.request(`${path}:cancel`, {});
}
response ??= await this.request(path);
response ??= await this.request(
waitForUnknown ? `${path}?reconcile=true` : path,
);
} catch (cause) {
if (
!waitForUnknown ||
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -261,7 +261,7 @@ private void execution(HttpExchange exchange, String suffix)
throw new RuntimeBrokerException(400, "runtime_payload_invalid", "payloadJson is required", false);
}
complete(exchange, service.startExecution(harnessSessionId, runtimeSessionId, executionCallId, payload)
.thenCompose(record -> observe(harnessSessionId, runtimeSessionId, record)),
.thenCompose(record -> observe(harnessSessionId, runtimeSessionId, record, false)),
observation -> observedExecutionEnvelope(harnessSessionId, runtimeSessionId, observation));
return;
}
Expand All @@ -278,7 +278,7 @@ private void execution(HttpExchange exchange, String suffix)
String runtimeSessionId = JsonCodec.requiredString(body,
"runtimeSessionId", "cancel request");
complete(exchange, service.cancelExecution(harnessSessionId, runtimeSessionId, executionCallId)
.thenCompose(record -> observe(harnessSessionId, runtimeSessionId, record)),
.thenCompose(record -> observe(harnessSessionId, runtimeSessionId, record, false)),
observation -> observedExecutionEnvelope(harnessSessionId, runtimeSessionId, observation));
return;
}
Expand All @@ -296,17 +296,24 @@ private void execution(HttpExchange exchange, String suffix)
if (query.containsKey("afterSeq")) {
parseSequence(query.get("afterSeq"));
}
String reconcile = query.getOrDefault("reconcile", "false");
if (!"true".equals(reconcile) && !"false".equals(reconcile)) {
throw new RuntimeBrokerException(400, "runtime_broker_invalid_request",
"reconcile must be true or false.", false);
}
complete(exchange, service.getExecution(harnessSessionId, runtimeSessionId, executionCallId)
.thenCompose(record -> observe(harnessSessionId, runtimeSessionId, record)),
.thenCompose(record -> observe(harnessSessionId, runtimeSessionId, record,
Boolean.parseBoolean(reconcile))),
observation -> observedExecutionEnvelope(harnessSessionId, runtimeSessionId, observation));
return;
}
throw notFound();
}

private CompletionStage<ExecutionReconciliation> observe(String harnessSessionId,
String runtimeSessionId, ToolExecutionRecord record) {
String runtimeSessionId, ToolExecutionRecord record, boolean reconcile) {
return record.getState() == ToolExecutionRecord.State.UNKNOWN
&& (reconcile || Integer.valueOf(3).equals(record.getReference().get("runtimeProtocol")))
? service.reconcileExecution(harnessSessionId, runtimeSessionId, record.getExecutionCallId())
: CompletableFuture.completedFuture(new ExecutionReconciliation(record,
ExecutionReconciliation.Outcome.IN_FLIGHT, null));
Expand Down
Loading
Loading