Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
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
17 changes: 17 additions & 0 deletions docs/design/2026-09-23-managed-runtime-process-adoption.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
# Managed Runtime process adoption

[English](2026-09-23-managed-runtime-process-adoption.md) | [简体中文](2026-09-23-managed-runtime-process-adoption.zh-CN.md)

Status: implemented. Updated: 2026-09-23. Continues the [attestation client](2026-09-23-java-runtime-attestation-client.md).

## This slice

The Broker starts a worker process, attests it, and only then stores the lease as READY. A later use of that in-memory lease attests again. If the process is gone, the call fails instead of reusing the old endpoint.
Comment thread
doudouOUC marked this conversation as resolved.
Outdated

The worker is the merged `managed-runtime-worker` command: one boot JSON document on stdin, one ready record on stdout. The preview `--boot-config` file launch is not used.

Tool HTTP (`POST /internal/managed-runtime/v2/execute`) is on the Java client. The merged worker still exposes only attestation, so execute against that process is a non-retryable 404. Mounting real tool handlers stays with the Hosted ordinary-tools slice.

## Not in this slice

Spring configuration and Flyway live with the Java control-plane module, which is not on `main`. Kubernetes provisioning stays out. This slice uses the existing in-memory and JDBC repositories; it does not add a server.
17 changes: 17 additions & 0 deletions docs/design/2026-09-23-managed-runtime-process-adoption.zh-CN.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
# Managed Runtime 进程接管

[English](2026-09-23-managed-runtime-process-adoption.md) | [简体中文](2026-09-23-managed-runtime-process-adoption.zh-CN.md)

状态:已实现。更新日期:2026-09-23。承接[attestation 客户端](2026-09-23-java-runtime-attestation-client.zh-CN.md)。

## 本切片

Broker 启动 worker 进程,证明通过之后才把 lease 记为 READY。之后再次使用这条内存中的 lease 时会重新证明。进程已经不在时,调用失败,不再复用旧 endpoint。

worker 使用已经合入的 `managed-runtime-worker`:标准输入一份 boot JSON,标准输出一条 ready 记录。不使用预览里的 `--boot-config` 文件启动。

Java 客户端提供工具 HTTP(`POST /internal/managed-runtime/v2/execute`)。已经合入的 worker 仍然只暴露 attestation,所以对这个进程执行工具会得到不可重试的 404。真正的工具处理留在 Hosted 普通工具那一笔。

## 不在本切片

Spring 配置和 Flyway 跟 Java 控制面模块走,那个模块还不在 `main` 上。Kubernetes provisioner 不包含在内。本切片使用现有的内存和 JDBC Repository,不新增服务器。
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,79 @@ public CompletionStage<RuntimeAttestation> attest(RuntimeLease lease,
return returned;
}

public CompletionStage<Void> execute(RuntimeLease lease,
Comment thread
doudouOUC marked this conversation as resolved.
RuntimeSession session, Map<String, Object> reference) {
Map<String, Object> body = sessionBody(session);
body.put("reference", reference);
return post(lease, "/internal/managed-runtime/v2/execute", body)
.thenApply(ignored -> null);
}

private Map<String, Object> sessionBody(RuntimeSession session) {
RuntimeScope scope = session.getScope();
Map<String, Object> body = new LinkedHashMap<>();
body.put("protocolVersion", 2);
body.put("tenantId", scope.getTenantId());
body.put("workspaceId", scope.getWorkspaceId());
body.put("workspaceCwd", scope.getCanonicalCwd());
body.put("sessionId", session.getRuntimeSessionId());
body.put("turnKind", session.getTurnKind());
return body;
}

private CompletionStage<byte[]> post(RuntimeLease lease, String path,
Map<String, Object> body) {
HttpRequest httpRequest = HttpRequest.newBuilder(
Comment thread
doudouOUC marked this conversation as resolved.
lease.getEndpoint().resolve(path))
.timeout(requestTimeout)
.header("Authorization", "Bearer " + lease.getToken())
.header("Cache-Control", "no-store")
.header("Content-Type", "application/json")
.header("X-Qwen-Managed-Lease-Id", lease.getLeaseId())
.header("X-Qwen-Managed-Lease-Epoch",
Long.toString(lease.getEpoch()))
.POST(HttpRequest.BodyPublishers.ofByteArray(
JsonCodec.encode(body)))
.build();
CompletableFuture<byte[]> result = new CompletableFuture<>();
CompletableFuture<HttpResponse<BoundedBody>> exchange = client
.sendAsync(httpRequest,
info -> new BoundedBodySubscriber(BODY_LIMIT_BYTES));
exchange.whenComplete((response, error) -> {
if (error != null) {
result.completeExceptionally(unavailable(unwrap(error)));
return;
}
BoundedBody responseBody = response.body();
if (responseBody.overflow() || response.statusCode() != 200) {
result.completeExceptionally(
Comment thread
doudouOUC marked this conversation as resolved.
failure(response.statusCode() == 200
? 413 : response.statusCode()));
return;
}
result.complete(responseBody.bytes());
});
CompletableFuture<byte[]> returned = result
.orTimeout(requestTimeout.toMillis(), TimeUnit.MILLISECONDS)
.handle((value, error) -> {
if (error == null) {
return value;
}
Throwable cause = unwrap(error);
if (cause instanceof RuntimeBrokerException failure) {
throw failure;
}
throw unavailable(cause);
});
returned.whenComplete((value, error) -> {
if (error != null || returned.isCancelled()) {
exchange.cancel(true);
result.cancel(false);
}
});
return returned;
}

private static Throwable unwrap(Throwable error) {
Throwable cause = error;
while (cause instanceof CompletionException
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,221 @@
package com.alibaba.qwen.code.runtimebroker;

import com.alibaba.fastjson2.JSONObject;
import java.io.BufferedReader;
import java.io.IOException;
import java.io.InputStreamReader;
import java.net.URI;
import java.nio.charset.StandardCharsets;
import java.nio.file.Path;
import java.security.SecureRandom;
import java.time.Duration;
import java.util.Base64;
import java.util.List;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionStage;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;

/**
* Starts the merged attestation worker over stdin and returns a lease only
* after that process attests as the same identity.
*/
public final class LocalProcessRuntimeProvisioner
Comment thread
doudouOUC marked this conversation as resolved.
implements RuntimeProvisioner, AutoCloseable {
private static final Duration READY_TIMEOUT = Duration.ofSeconds(30);
private static final SecureRandom RANDOM = new SecureRandom();

private final List<String> command;
private final Path workingDirectory;
private final HttpRuntimeTransport transport;
private final ExecutorService executor = Executors.newCachedThreadPool(
task -> {
Thread thread = new Thread(task, "runtime-provisioner");
thread.setDaemon(true);
return thread;
});
private final ConcurrentMap<String, OwnedProcess> owned =
new ConcurrentHashMap<>();

public LocalProcessRuntimeProvisioner(List<String> command,
Path workingDirectory, HttpRuntimeTransport transport) {
if (command == null || command.isEmpty() || workingDirectory == null
|| transport == null) {
throw new IllegalArgumentException(
"worker command, directory, and transport are required");
}
this.command = List.copyOf(command);
this.workingDirectory = workingDirectory;
this.transport = transport;
}

@Override
public CompletionStage<RuntimeLease> provision(
RuntimeProvisionRequest request) {
return CompletableFuture.supplyAsync(() -> start(request), executor);
}

@Override
public CompletionStage<Void> confirm(RuntimeProvisionRequest request,
RuntimeLease lease) {
return CompletableFuture.runAsync(() -> attestOwned(request, lease),
executor);
}

void stop(RuntimeLease lease) {
Comment thread
doudouOUC marked this conversation as resolved.
OwnedProcess process = owned.remove(lease.getRuntimeInstanceId());
Comment thread
doudouOUC marked this conversation as resolved.
if (process != null) {
process.process.destroy();
}
}

@Override
public void close() {
Comment thread
doudouOUC marked this conversation as resolved.
for (OwnedProcess process : owned.values()) {
process.process.destroy();
}
owned.clear();
executor.shutdownNow();
}

private RuntimeLease start(RuntimeProvisionRequest request) {
OwnedProcess ownedProcess = null;
try {
String runtimeInstanceId = UUID.randomUUID().toString();
String runtimeIncarnation = UUID.randomUUID().toString();
String leaseId = UUID.randomUUID().toString();
String provisionRequestId = UUID.randomUUID().toString();
byte[] tokenBytes = new byte[32];
RANDOM.nextBytes(tokenBytes);
String token = Base64.getUrlEncoder().withoutPadding()
.encodeToString(tokenBytes);
RuntimeScope scope = request.getScope();
JSONObject boot = new JSONObject();
boot.put("capabilityDigest", scope.getCapabilityDigest());
Comment thread
doudouOUC marked this conversation as resolved.
boot.put("epoch", 1);
boot.put("isolationClass", scope.getIsolationClass());
boot.put("leaseId", leaseId);
boot.put("provisionRequestId", provisionRequestId);
boot.put("runtimeIncarnation", runtimeIncarnation);
boot.put("runtimeInstanceId", runtimeInstanceId);
boot.put("tenantId", scope.getTenantId());
boot.put("token", token);
boot.put("type", "boot");
boot.put("version", 1);
boot.put("workspaceCwd", scope.getCanonicalCwd());
boot.put("workspaceGeneration", scope.getWorkspaceGeneration());
boot.put("workspaceId", scope.getWorkspaceId());
Process process = new ProcessBuilder(command)
Comment thread
doudouOUC marked this conversation as resolved.
.directory(workingDirectory.toFile())
.redirectError(ProcessBuilder.Redirect.DISCARD)
Comment thread
doudouOUC marked this conversation as resolved.
.start();
ownedProcess = new OwnedProcess(process, request,
new RuntimeProvisionSeed(provisionRequestId,
runtimeInstanceId, runtimeIncarnation, leaseId,
1, token));
process.getOutputStream().write(boot.toJSONString()
.getBytes(StandardCharsets.UTF_8));
process.getOutputStream().close();
String readyLine = readReadyLine(process);
Map<String, Object> ready = JsonCodec.parseObject(
readyLine.getBytes(StandardCharsets.UTF_8),
"Managed Runtime ready record");
if (!"ready".equals(ready.get("type"))
|| !Long.valueOf(1L).equals(number(ready.get("version")))
|| !runtimeInstanceId.equals(
ready.get("runtimeInstanceId"))
|| !runtimeIncarnation.equals(
ready.get("runtimeIncarnation"))
|| !leaseId.equals(ready.get("leaseId"))
|| !Long.valueOf(1L).equals(number(ready.get("epoch")))) {
throw failed("Managed Runtime ready record is invalid.");
}
RuntimeLease lease = new RuntimeLease(runtimeInstanceId,
URI.create(String.valueOf(ready.get("url"))), token,
leaseId, 1);
attest(request, ownedProcess.seed, lease);
owned.put(runtimeInstanceId, ownedProcess);
return lease;
} catch (RuntimeException exception) {
Comment thread
doudouOUC marked this conversation as resolved.
Outdated
if (ownedProcess != null) {
ownedProcess.process.destroy();
}
throw exception;
} catch (IOException exception) {
if (ownedProcess != null) {
ownedProcess.process.destroy();
}
throw failed("Managed Runtime worker failed to start.",
exception);
}
}

private void attestOwned(RuntimeProvisionRequest request,
RuntimeLease lease) {
OwnedProcess process = owned.get(lease.getRuntimeInstanceId());
if (process == null || !process.process.isAlive()) {
throw failed("Managed Runtime process is not alive.");
}
attest(request, process.seed, lease);
}

private void attest(RuntimeProvisionRequest request,
RuntimeProvisionSeed seed, RuntimeLease lease) {
try {
transport.attest(lease, request, seed).toCompletableFuture()
.get(READY_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS);
} catch (Exception exception) {
throw failed("Managed Runtime attestation failed.", exception);
}
}

private static String readReadyLine(Process process) throws IOException {
CompletableFuture<String> line = new CompletableFuture<>();
Thread reader = new Thread(() -> {
try (BufferedReader input = new BufferedReader(
Comment thread
doudouOUC marked this conversation as resolved.
Outdated
Comment thread
doudouOUC marked this conversation as resolved.
Outdated
new InputStreamReader(process.getInputStream(),
StandardCharsets.UTF_8))) {
line.complete(input.readLine());
} catch (IOException exception) {
Comment thread
doudouOUC marked this conversation as resolved.
Outdated
line.completeExceptionally(exception);
}
}, "runtime-ready");
reader.setDaemon(true);
reader.start();
try {
String ready = line.get(READY_TIMEOUT.toMillis(),
TimeUnit.MILLISECONDS);
if (ready == null) {
throw failed("Managed Runtime worker closed before ready.");
Comment thread
doudouOUC marked this conversation as resolved.
Outdated
}
return ready;
} catch (Exception exception) {
process.destroyForcibly();
throw failed("Managed Runtime worker did not become ready.",
exception);
}
}

private static Long number(Object value) {
return value instanceof Number number ? number.longValue() : null;
}

private static RuntimeBrokerException failed(String message) {
return failed(message, null);
}

private static RuntimeBrokerException failed(String message,
Throwable cause) {
return new RuntimeBrokerException(503, "runtime_provision_failed",
message, true, cause);
}

private record OwnedProcess(Process process,
Comment thread
doudouOUC marked this conversation as resolved.
Outdated
RuntimeProvisionRequest request, RuntimeProvisionSeed seed) {
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -486,8 +486,9 @@ private CompletionStage<BindingContext> ensureBinding(
if (finishing != null) {
return finishing;
}
return CompletableFuture.completedFuture(
requireLiveBinding(record));
BindingContext live = requireLiveBinding(record);
Comment thread
doudouOUC marked this conversation as resolved.
Comment thread
doudouOUC marked this conversation as resolved.
return safeStage(() -> provisioner.confirm(request, live.lease()))
Comment thread
doudouOUC marked this conversation as resolved.
Outdated
.thenApply(ignored -> live);
}
if (record.getState()
!= RuntimeBindingRecord.State.PROVISIONING) {
Expand Down
Original file line number Diff line number Diff line change
@@ -1,9 +1,19 @@
package com.alibaba.qwen.code.runtimebroker;

import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionStage;

/** Ensures one physical Runtime resource for a claimed placement. */
public interface RuntimeProvisioner {
/** Retries for the same request must converge on one live resource. */
CompletionStage<RuntimeLease> provision(RuntimeProvisionRequest request);

/**
* Proves a lease this process already treats as ready still answers
* attestation. The default accepts the in-memory lease.
*/
default CompletionStage<Void> confirm(RuntimeProvisionRequest request,
RuntimeLease lease) {
return CompletableFuture.completedFuture(null);
}
}
Loading
Loading