diff --git a/apps/daemon/internal/agent/claudesdk/executor.go b/apps/daemon/internal/agent/claudesdk/executor.go index cf53eb6e1..ebbb1f6f2 100644 --- a/apps/daemon/internal/agent/claudesdk/executor.go +++ b/apps/daemon/internal/agent/claudesdk/executor.go @@ -35,7 +35,6 @@ func startExecutor(ctx context.Context, checked *runtimeCheckCache, probe Config if err = validateExecutorFeatures(info, start); err != nil { return nil, err } - start.Type = "executor_prepare" base, err := run() if err != nil { return nil, err diff --git a/apps/daemon/internal/agent/claudesdk/executor_fixture_test.go b/apps/daemon/internal/agent/claudesdk/executor_fixture_test.go index 5f2241fa4..3ee557c67 100644 --- a/apps/daemon/internal/agent/claudesdk/executor_fixture_test.go +++ b/apps/daemon/internal/agent/claudesdk/executor_fixture_test.go @@ -8,6 +8,7 @@ import ( "encoding/json" "os" "path/filepath" + "time" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent/clirunner" @@ -59,7 +60,7 @@ func startSingleTurn(ctx context.Context, config testBridge, req proto.PromptReq if err != nil { return nil, err } - resource, err := config.factory()(ctx, agent.PrepareRequest{PromptRequestPayload: req, Prepared: configuration}) + resource, err := config.factory()(ctx, agent.PrepareRequest{PreparationDeadline: time.Now().Add(time.Minute), PromptRequestPayload: req, Prepared: configuration}) if err != nil { return nil, err } diff --git a/apps/daemon/internal/agent/claudesdk/options.go b/apps/daemon/internal/agent/claudesdk/options.go index 81106aeba..13387fefc 100644 --- a/apps/daemon/internal/agent/claudesdk/options.go +++ b/apps/daemon/internal/agent/claudesdk/options.go @@ -1,7 +1,9 @@ package claudesdk import ( + "context" "fmt" + "time" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" @@ -21,26 +23,30 @@ type subagentOptions struct { } type startRequest struct { - NativeModelOptions *nativeModelOptions `json:"native_model_options,omitempty"` - ToolSearch bool `json:"tool_search,omitempty"` - Subagents *subagentOptions `json:"subagents,omitempty"` - OutputFormat *proto.OutputFormat `json:"output_format,omitempty"` - Type string `json:"type"` - Model string `json:"model"` - SystemPrompt string `json:"system_prompt"` - Cwd string `json:"cwd"` - Resume string `json:"resume,omitempty"` - Functions []proto.FunctionTool `json:"functions,omitempty"` - MCPHTTPServers *[]mcpHTTPServer `json:"mcp_http_servers,omitempty"` - Workspace *workspaceProfile `json:"workspace,omitempty"` - RequireHistory bool `json:"require_history,omitempty"` + PreparationDeadline int64 `json:"preparation_deadline"` + NativeModelOptions *nativeModelOptions `json:"native_model_options,omitempty"` + ToolSearch bool `json:"tool_search,omitempty"` + Subagents *subagentOptions `json:"subagents,omitempty"` + OutputFormat *proto.OutputFormat `json:"output_format,omitempty"` + Type string `json:"type"` + Model string `json:"model"` + SystemPrompt string `json:"system_prompt"` + Cwd string `json:"cwd"` + Resume string `json:"resume,omitempty"` + Functions []proto.FunctionTool `json:"functions,omitempty"` + MCPHTTPServers *[]mcpHTTPServer `json:"mcp_http_servers,omitempty"` + Workspace *workspaceProfile `json:"workspace,omitempty"` + RequireHistory bool `json:"require_history,omitempty"` } // prepareOptions renders the request's execution configuration and the // selected model provider. The registered factory already admitted the // selection against the declaration. func prepareOptions(req agent.PrepareRequest) (startRequest, []string, error) { - start := startRequest{Type: "start", Resume: req.AgentSessionID, RequireHistory: req.RequireExistingNativeSession, ToolSearch: req.ToolSearch} + if req.PreparationDeadline.IsZero() || !time.Now().Before(req.PreparationDeadline) { + return startRequest{}, nil, context.DeadlineExceeded + } + start := startRequest{PreparationDeadline: req.PreparationDeadline.UnixMilli(), Type: "executor_prepare", Resume: req.AgentSessionID, RequireHistory: req.RequireExistingNativeSession, ToolSearch: req.ToolSearch} fail := func(reason string) (startRequest, []string, error) { return startRequest{}, nil, fmt.Errorf("claudesdk: %s", reason) } diff --git a/apps/daemon/internal/agent/claudesdk/options_test.go b/apps/daemon/internal/agent/claudesdk/options_test.go index d1103eb3d..50f76285d 100644 --- a/apps/daemon/internal/agent/claudesdk/options_test.go +++ b/apps/daemon/internal/agent/claudesdk/options_test.go @@ -1,9 +1,11 @@ package claudesdk import ( + "context" "errors" "slices" "testing" + "time" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent/clirunner" @@ -55,10 +57,26 @@ func prepared(t testing.TB, req proto.PromptRequestPayload) agent.PrepareRequest if err != nil { t.Fatal(err) } - return agent.PrepareRequest{PromptRequestPayload: req, Prepared: configuration} + return agent.PrepareRequest{PreparationDeadline: time.Now().Add(time.Minute), PromptRequestPayload: req, Prepared: configuration} } // fixtureProvider is the provider every Claude request carries. func fixtureProvider() *modelprovider.Provider { return &modelprovider.Provider{Protocol: modelprovider.Anthropic, BaseURL: "https://model.example", APIKey: "fixture-key"} } + +func TestPreparationDeadlineReachesBridgeUnchanged(t *testing.T) { + req := prepared(t, proto.PromptRequestPayload{ModelProvider: fixtureProvider(), DisableExecutionEnvironment: true, Model: "test-model"}) + deadline := time.Now().Add(3 * time.Minute).Truncate(time.Millisecond) + req.PreparationDeadline = deadline + start, _, err := prepareTestView(t, req) + if err != nil || start.PreparationDeadline != deadline.UnixMilli() { + t.Fatalf("deadline=%d error=%v", start.PreparationDeadline, err) + } + for _, invalid := range []time.Time{{}, time.Now().Add(-time.Second)} { + req.PreparationDeadline = invalid + if _, _, err := prepareTestView(t, req); !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("invalid budget: %v", err) + } + } +} diff --git a/apps/daemon/internal/agent/harness.go b/apps/daemon/internal/agent/harness.go index 1f810f106..eebd4a70c 100644 --- a/apps/daemon/internal/agent/harness.go +++ b/apps/daemon/internal/agent/harness.go @@ -41,6 +41,7 @@ import ( "slices" "strconv" "strings" + "time" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent/clirunner" "github.com/MiniMax-AI/OpenAgentCore/internal/agentcapabilities" @@ -532,6 +533,10 @@ type ExecutorFactory func(context.Context, PrepareRequest) (Executor, error) // sets Prepared. A Turn's run ID and input arrive in Executor.StartTurn. type PrepareRequest struct { proto.PromptRequestPayload + // PreparationDeadline is the dispatch-owned absolute preparation deadline. + // It must be nonzero and preserved through preparation, never restarted. + // The factory context owns the Executor lifetime; this deadline does not. + PreparationDeadline time.Time // Prepared is the model configuration, validated against the kind's // declaration. An adapter takes its model, provider and native parameters // only from here. diff --git a/apps/daemon/internal/agenthost/agenthost_linux_test.go b/apps/daemon/internal/agenthost/agenthost_linux_test.go index ae869e6c9..48b696cf1 100644 --- a/apps/daemon/internal/agenthost/agenthost_linux_test.go +++ b/apps/daemon/internal/agenthost/agenthost_linux_test.go @@ -174,6 +174,9 @@ func newDaemon(t *testing.T, cfg Config, d deps) *daemon { } e, err := dm.host.openExecutor(ctx, req) if s, ok := e.(*session); ok { + if req.PreparationDeadline.IsZero() || s.plan.request.PreparationDeadline != req.PreparationDeadline { + t.Error("agent host changed the dispatch preparation deadline") + } dm.mu.Lock() dm.opened[req.Assignment.SessionID] = s dm.mu.Unlock() diff --git a/apps/daemon/internal/dispatch/executor.go b/apps/daemon/internal/dispatch/executor.go index 9384be70e..7dfea527b 100644 --- a/apps/daemon/internal/dispatch/executor.go +++ b/apps/daemon/internal/dispatch/executor.go @@ -186,7 +186,7 @@ func (r *Router) prepareExecutor(p *preparationState, configuration proto.Prompt } var native agent.Executor var err error - req := agent.PrepareRequest{PromptRequestPayload: configuration, StateKey: "agents-api-" + owner.sessionID, Assignment: p.request.Assignment} + req := agent.PrepareRequest{PromptRequestPayload: configuration, StateKey: "agents-api-" + owner.sessionID, Assignment: p.request.Assignment, PreparationDeadline: p.deadline} if owner.ctx.Err() == nil { if environment != nil { req, err = environment.Prepare(owner.ctx, req) diff --git a/apps/daemon/internal/dispatch/executor_test.go b/apps/daemon/internal/dispatch/executor_test.go index f06e97663..95bb8cea5 100644 --- a/apps/daemon/internal/dispatch/executor_test.go +++ b/apps/daemon/internal/dispatch/executor_test.go @@ -347,3 +347,68 @@ func (*reusableTurn) SteerWithReceipt(context.Context, proto.PromptSteerPayload, func (*reusableTurn) SubmitFunctionResult(context.Context, proto.FunctionResultPayload) error { return agent.ErrUnknownFunctionCall } + +func TestExecutorPreparationDeadlineIsPreservedWithoutOwningWarmLifetime(t *testing.T) { + environment := newTestOwner(preparationEnvironmentID, preparationSessionID) + beforeEnvironment, afterEnvironment := make(chan agent.PrepareRequest, 1), make(chan agent.PrepareRequest, 1) + proceed := make(chan struct{}) + environment.prepare = func(req agent.PrepareRequest) (agent.PrepareRequest, error) { + beforeEnvironment <- req + <-proceed + req.WorkspaceRoot = "/workspace" + return req, nil + } + owner := &reusableExecutor{starts: make(chan *reusableTurn, 2)} + ownerContext := make(chan context.Context, 1) + reg := agent.NewRegistry() + reg.RegisterKind(proto.SupportedAgentKind{Kind: "prepared", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{LocalEnvironment: proto.CapabilitySupported, FunctionTools: proto.CapabilitySupported, FunctionResultImages: proto.CapabilitySupported, NativeSessionRecovery: proto.CapabilitySupported})}, prototest.ModelConfiguration()) + reg.RegisterExecutor("prepared", func(ctx context.Context, req agent.PrepareRequest) (agent.Executor, error) { + ownerContext <- ctx + afterEnvironment <- req + return owner, nil + }) + sender := &recSender{} + r, err := dispatch.New(dispatch.Config{Registry: reg, Sender: sender, Environments: environment.Resolve, PreparationTimeout: time.Second, IdleTimeout: time.Minute}) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + if err := r.Shutdown(ctx); err != nil { + t.Error(err) + } + }) + assign(t, r, preparationSessionID, preparationEnvironmentID) + req := preparationRequest() + if err := r.Handle(t.Context(), mustEnv(t, proto.TypeExecutionPrepare, "deadline", req)); err != nil { + t.Fatal(err) + } + admitted := waitPreparationStatus(t, sender, "deadline", "preparing", "") + before := <-beforeEnvironment + close(proceed) + after := <-afterEnvironment + ctx := <-ownerContext + if before.PreparationDeadline.IsZero() || before.PreparationDeadline != after.PreparationDeadline || before.PreparationDeadline.UnixMilli() != admitted.ExpiresAt || after.WorkspaceRoot != "/workspace" || after.Prepared.Model == "" { + t.Fatalf("preparation handoff changed: before=%v after=%v admission=%d", before.PreparationDeadline, after.PreparationDeadline, admitted.ExpiresAt) + } + if _, hasDeadline := ctx.Deadline(); hasDeadline { + t.Fatal("preparation deadline attached to Executor context") + } + ready := waitPreparationStatus(t, sender, "deadline", "ready", "") + startExecutorTurn(t, r, sender, "deadline", "first", ready) + first := <-owner.starts + <-time.After(time.Until(before.PreparationDeadline)) + if ctx.Err() != nil || owner.closes.Load() != 0 { + t.Fatal("adopted Executor ended at preparation deadline") + } + first.finish() + waitFor(t, func() bool { return r.ActiveRuns() == 0 }, "first settled") + req.Configuration.AgentSessionID = "native-session" + next := executorAdmission(t, r, sender, "next", req) + if !next.Reused || next.ExecutorID != ready.ExecutorID { + t.Fatal("warm owner not reused after deadline") + } + startExecutorTurn(t, r, sender, "next", "second", next) + (<-owner.starts).finish() +} diff --git a/contracts/agents-api/harness-onboarding.md b/contracts/agents-api/harness-onboarding.md index 17432541a..5251b7ce1 100644 --- a/contracts/agents-api/harness-onboarding.md +++ b/contracts/agents-api/harness-onboarding.md @@ -59,7 +59,7 @@ Implement the mandatory text lifecycle and handle every extension explicitly. Qu [`agent/harness.go`](https://github.com/MiniMax-AI/OpenAgentCore/blob/main/apps/daemon/internal/agent/harness.go) is the interface entry point. The required lifecycle is `ViewExecutorFactory`, `Executor`, `Turn` and `TurnSettlement`. `Turn` is one interface: `Cancel`, `CancellationOutcome`, `SteerWithReceipt`, `SubmitFunctionResult` and `AwaitSettlement`. Required methods perform their native obligations; returning Unsupported is not an implementation of cancellation, receipts, settlement or cleanup. An operation the adapter does not support returns Unsupported, and the capability declaration, not the method, decides whether the Runtime calls it. All use the neutral protocol types. -The view's `ViewExecutorFactory` takes one `agent.PrepareRequest`: the Session's configuration as `execution_prepare` carries it, the model configuration that the Registry prepared once from the kind's declaration (`Prepared`), the Session's native state key (`StateKey`), and the Environment's workspace and installed Capabilities (`WorkspaceRoot`, `CapabilityRoot`, `Skills`, `MCP`), which its owner fills. The adapter takes its model, provider and native parameters only from `Prepared` and never parses `model` or `model_provider` itself. A Turn's Run ID and input arrive in `Executor.StartTurn`. +The view's `ViewExecutorFactory` takes one `agent.PrepareRequest`: the Session's configuration as `execution_prepare` carries it, the model configuration that the Registry prepared once from the kind's declaration (`Prepared`), the Session's native state key (`StateKey`), and the Environment's workspace and installed Capabilities (`WorkspaceRoot`, `CapabilityRoot`, `Skills`, `MCP`), which its owner fills. The adapter takes its model, provider and native parameters only from `Prepared` and never parses `model` or `model_provider` itself. Dispatch also supplies the required absolute `PreparationDeadline`, created once for the admission. Environment preparation, Registry and View preserve it unchanged. Native initialization must finish before it; an adapter that passes an initialization budget to its SDK derives it from the remaining time and rejects an exhausted budget. The factory context owns the Executor lifetime independently: adopting a prepared Executor must not attach the preparation deadline to its later Turns. A Turn's Run ID and input arrive in `Executor.StartTurn`. For example, the Codex adapter keeps its app-server and thread, the Claude adapter one streaming Query, and the MiniMax adapter its ACP connection and native session. All expose the same Executor and Turn contract. Native callbacks and resources stay inside the adapter; the Runtime owns admission, idle expiry and replacement. Cancellation targets the exact Turn through `Turn.Cancel`, and the adapter supplies native completion evidence to the Runtime. diff --git a/contracts/agents-api/zh/harness-onboarding.md b/contracts/agents-api/zh/harness-onboarding.md index d274d9240..005bd3e66 100644 --- a/contracts/agents-api/zh/harness-onboarding.md +++ b/contracts/agents-api/zh/harness-onboarding.md @@ -1,7 +1,7 @@ --- title: "添加 Harness" source: contracts/agents-api/harness-onboarding.md -source_hash: 56b2a40b814062de7d636af33ad2eb7671b8b3eb11781e32caacdf3d3985ae73 +source_hash: a18a2221b28aa01596574a1db63e1bf643dc43fb73c15261fb8932dd447f499b --- **Harness** 是一种运行模型和工具循环的原生代理引擎(Codex、Claude Code、MiniMax Code)。**Harness 适配器**将 Runtime 的 Executor 和 Turn 契约转换到该引擎的 SDK 或协议。本文档定义 Runtime–Harness 协议:适配器接口及其生命周期义务、注册、支持声明和验收。 @@ -61,7 +61,7 @@ Environment 提供执行资源。受管 E2B、Docker 和 microsandbox 机器以 [`agent/harness.go`](https://github.com/MiniMax-AI/OpenAgentCore/blob/main/apps/daemon/internal/agent/harness.go) 是接口入口。必需的生命周期包括 `ViewExecutorFactory`、`Executor`、`Turn` 和 `TurnSettlement`。`Turn` 是一个接口:`Cancel`、`CancellationOutcome`、`SteerWithReceipt`、`SubmitFunctionResult` 和 `AwaitSettlement`。必需方法必须履行其原生义务;返回 Unsupported 并不构成对取消、回执、结算或清理的实现。适配器不支持的操作返回 Unsupported,由能力声明而不是方法决定 Runtime 是否调用它。所有接口都使用中立协议类型。 -视图的 `ViewExecutorFactory` 接收一个 `agent.PrepareRequest`:`execution_prepare` 携带的 Session 配置、Registry 按 kind 的声明一次性准备好的模型配置(`Prepared`)、Session 的原生状态键(`StateKey`),以及由 Environment owner 填写的 Environment 工作区和已安装 Capabilities(`WorkspaceRoot`、`CapabilityRoot`、`Skills`、`MCP`)。适配器只从 `Prepared` 获取模型、提供商和原生参数,从不自行解析 `model` 或 `model_provider`。Turn 的 Run ID 和输入通过 `Executor.StartTurn` 传入。 +视图的 `ViewExecutorFactory` 接收一个 `agent.PrepareRequest`:`execution_prepare` 携带的 Session 配置、Registry 按 kind 的声明一次性准备好的模型配置(`Prepared`)、Session 的原生状态键(`StateKey`),以及由 Environment owner 填写的 Environment 工作区和已安装 Capabilities(`WorkspaceRoot`、`CapabilityRoot`、`Skills`、`MCP`)。适配器只从 `Prepared` 获取模型、提供商和原生参数,从不自行解析 `model` 或 `model_provider`。Dispatch 还提供 admission 创建时唯一确定的必填绝对截止时间 `PreparationDeadline`。Environment 准备、Registry 与 View 原样传递它。原生初始化必须在截止时间前结束;适配器向 SDK 传递初始化预算时,必须从剩余时间推导,并拒绝已耗尽的预算。工厂上下文独立拥有 Executor 生命周期:接管准备完成的 Executor 后,不得让准备截止时间约束后续 Turn。Turn 的 Run ID 和输入通过 `Executor.StartTurn` 传入。 例如,Codex 适配器保留其 app-server 和 thread,Claude 适配器保留一个流式 Query,MiniMax 适配器保留其 ACP 连接和原生 session。它们都公开相同的 Executor 和 Turn 契约。原生回调和资源保留在适配器内部;Runtime 负责准入、空闲过期和替换。取消通过 `Turn.Cancel` 精确定位到目标 Turn,适配器则向 Runtime 提供原生完成证据。 diff --git a/packages/claude-sdk-adapter/README.md b/packages/claude-sdk-adapter/README.md index aeaf289e4..8e532dbad 100644 --- a/packages/claude-sdk-adapter/README.md +++ b/packages/claude-sdk-adapter/README.md @@ -34,7 +34,7 @@ The Harness environment is closed: the Session home's native directories, the ga ### Executor preparation and Turns -The private bridge accepts `executor_prepare` without model input. It freezes validated configuration and resume identity, checks required history, and retains one native process and SDK Query across Turns. Preparation requires initialization and acknowledgement of required hooks while the input iterator remains empty. The pinned SDK owns the initialization deadline through its default `startup` behavior. An `executor_ready` receipt permits later `turn_start` messages containing only Turn identity and ordered input; configuration replacement and concurrent starts are rejected. Every Turn event carries its originating `turn_id`. Native Session identity and actual tool inventory are checked before `input_ready`. +The private bridge accepts `executor_prepare` without model input. It freezes validated configuration and resume identity, checks required history, and retains one native process and SDK Query across Turns. Preparation requires initialization and acknowledgement of required hooks while the input iterator remains empty. The required `preparation_deadline` is dispatch’s absolute Unix millisecond deadline from `PrepareRequest.PreparationDeadline`. Immediately before `startup`, the bridge passes the remaining milliseconds to the pinned SDK’s `initializeTimeoutMs`; an exhausted budget fails without starting native initialization. Readiness must still precede that deadline. The budget applies only to preparation and does not terminate an adopted warm Query. An `executor_ready` receipt permits later `turn_start` messages containing only Turn identity and ordered input; configuration replacement and concurrent starts are rejected. Every Turn event carries its originating `turn_id`. Native Session identity and actual tool inventory are checked before `input_ready`. Each Turn ends with a result/error and `turn_settled`, independently of process exit. The outer input iterator remains open for later Turns. `turn_cancel` invokes the native interrupt control for that exact Turn. Unconfirmed input, native child work or queue state invalidates the Executor and requires close before replacement. EOF, owner signals and invalid control input close owned resources. Preparation may write native metadata and perform startup traffic; readiness does not prove provider authentication, complete sandbox health or placement authorization. diff --git a/packages/claude-sdk-adapter/src/adapter.ts b/packages/claude-sdk-adapter/src/adapter.ts index 5a850374b..485887555 100644 --- a/packages/claude-sdk-adapter/src/adapter.ts +++ b/packages/claude-sdk-adapter/src/adapter.ts @@ -125,7 +125,9 @@ export async function execute(request: Start | Prepare | ExecutorPrepare, emit: if(canUseTool) options.canUseTool=(...args)=>turns.track(()=>canUseTool(...args)); } if (request.type === "prepare" || request.type === "executor_prepare") { - warm = await startup({ options }); + const initializeTimeoutMs = request.type === "executor_prepare" ? request.preparation_deadline - Date.now() : undefined; + if (initializeTimeoutMs !== undefined && (!Number.isSafeInteger(initializeTimeoutMs) || initializeTimeoutMs <= 0)) throw new Error("preparation expired"); + warm = await startup({ options, initializeTimeoutMs }); stream = warm.query(turns ?? inputs); const initialized = await stream.initializationResult(); const hooksRequested = Object.values(options.hooks ?? {}).some(matchers => matchers?.some(matcher => matcher.hooks.length > 0)); @@ -147,6 +149,7 @@ export async function execute(request: Start | Prepare | ExecutorPrepare, emit: if(value.type==="steer") for(const event of inputs.submit(value)) await emit(event); else functions.submit(JSON.stringify(value)); }); + if (request.type !== "executor_prepare" || Date.now() >= request.preparation_deadline) throw new Error("preparation expired"); await ownerEmit({type:"executor_ready",protocol:3}); } else await emit({ type: "prepared" }); } else if (profile) { diff --git a/packages/claude-sdk-adapter/src/request.ts b/packages/claude-sdk-adapter/src/request.ts index ceedb5b3e..d2d5eb1cf 100644 --- a/packages/claude-sdk-adapter/src/request.ts +++ b/packages/claude-sdk-adapter/src/request.ts @@ -22,7 +22,7 @@ export type Start = { mcp_http_servers?: HTTPServer[]; workspace?: Workspace; }; -export type ExecutorPrepare = Omit & { type: "executor_prepare" }; +export type ExecutorPrepare = Omit & { type: "executor_prepare"; preparation_deadline: number }; export type Prepare = Omit & { type: "prepare"; workspace: Workspace }; // MCP startup confirms its hooks before the native input iterator yields. @@ -34,7 +34,7 @@ export function parseRequest(line: string): Start | Prepare | ExecutorPrepare { const value: unknown = JSON.parse(line); if (!value || typeof value !== "object" || Array.isArray(value)) throw new Error("invalid_request"); const request = value as Record; - const allowed = new Set(["native_model_options", "type", "input", "model", "system_prompt", "cwd", "resume", "require_history", "output_format", "subagents", "functions", "tool_search", "mcp_http_servers", "workspace"]); + const allowed = new Set(["preparation_deadline", "native_model_options", "type", "input", "model", "system_prompt", "cwd", "resume", "require_history", "output_format", "subagents", "functions", "tool_search", "mcp_http_servers", "workspace"]); if (Object.keys(request).some(key => !allowed.has(key)) || (request.type !== "start" && request.type !== "prepare" && request.type !== "executor_prepare") || (request.type === "start" ? !Array.isArray(request.input) : "input" in request) || @@ -43,6 +43,9 @@ export function parseRequest(line: string): Start | Prepare | ExecutorPrepare { typeof request.cwd !== "string" || !isAbsolute(request.cwd) || (request.require_history !== undefined && typeof request.require_history !== "boolean") || (request.resume !== undefined && (typeof request.resume !== "string" || !request.resume))) throw new Error("invalid_request"); + if (request.type === "executor_prepare" ? + !Number.isSafeInteger(request.preparation_deadline) || (request.preparation_deadline as number) <= 0 : + "preparation_deadline" in request) throw new Error("invalid_request"); if (request.functions !== undefined && (!Array.isArray(request.functions) || request.functions.some(tool => !tool || typeof tool.name !== "string" || !tool.name || typeof tool.description !== "string" || (tool.defer_loading !== undefined && typeof tool.defer_loading !== "boolean") || diff --git a/packages/claude-sdk-adapter/tests/executor.test.mjs b/packages/claude-sdk-adapter/tests/executor.test.mjs index a94de9510..a9fe08eda 100644 --- a/packages/claude-sdk-adapter/tests/executor.test.mjs +++ b/packages/claude-sdk-adapter/tests/executor.test.mjs @@ -23,7 +23,7 @@ globalThis.startupFixture=async({options})=>{ const close=()=>{child.stdin.end();interrupt?.()}; options.abortController.signal.addEventListener('abort',close); return {close,query(prompt){assert.equal(queried,false);queried=true;process.send({type:'query'});return { - close,async initializationResult(){return process.argv[1]==="no-hook-report" ? {} : process.argv[1]==="false-hook-report" ? {hooks_applied:false} : {hooks_applied:true}}, + close,async initializationResult(){if(process.argv[1]==="late-ready"){const expired=Date.now()+120000;Date.now=()=>expired;}return process.argv[1]==="no-hook-report" ? {} : process.argv[1]==="false-hook-report" ? {hooks_applied:false} : {hooks_applied:true}}, async interrupt(){process.send({type:'interrupt'});interrupt?.();if(process.argv[1]==='rejected')throw new Error('interrupt failed');return ['unknown','pending-function-unknown'].includes(process.argv[1]) ? undefined : {still_queued:process.argv[1]==='queued' ? ['not-consumed'] : []}}, async *[Symbol.asyncIterator](){ for await(const user of prompt){ @@ -77,9 +77,9 @@ async function launch(t,mode="normal") { const wait=async predicate=>{const until=Date.now()+5000;while(!predicate()){assert.ok(Date.now()setTimeout(resolve,5))}}; const send=value=>child.stdin.write(JSON.stringify(value)+"\n"); const start=(id,text)=>send({type:"turn_start",turn_id:id,input:[{content:[{type:"input_text",text}]}]}); - send({type:"executor_prepare",cwd:"/tmp",model:"fixture",system_prompt:"", + send({type:"executor_prepare",preparation_deadline:Date.now()+60000,cwd:"/tmp",model:"fixture",system_prompt:"", ...(mode==="features" || mode.startsWith("pending-function") ? {functions:[{name:"lookup",description:"lookup",parameters:{type:"object",properties:{text:{type:"string"}}}}]} : {})}); - await wait(()=>events.some(event=>event.type==="executor_ready")); + await wait(()=>events.some(event=>event.type===(mode==="late-ready"?"error":"executor_ready"))); assert.equal(observations.filter(event=>event.type==="input").length,0); return {child,events,observations,closed,wait,send,start}; } @@ -193,3 +193,11 @@ for (const mode of ["pending-function-result", "pending-function-terminal", "pen } assert.deepEqual(await closed,{code:0,signal:null}); }); + +test("an initialization finishing past the original deadline never publishes ready or consumes input",{timeout:10000},async t=>{ + const {events,observations,closed}=await launch(t,"late-ready"); + assert.deepEqual(await closed,{code:0,signal:null}); + assert.equal(events.some(event=>event.type==="executor_ready"),false); + assert.equal(observations.some(event=>event.type==="input"),false); + assert.equal(observations.filter(event=>event.type==="native_closed").length,1); +}); diff --git a/packages/claude-sdk-adapter/tests/native_model_options.test.mjs b/packages/claude-sdk-adapter/tests/native_model_options.test.mjs index 0759115da..1956ee57c 100644 --- a/packages/claude-sdk-adapter/tests/native_model_options.test.mjs +++ b/packages/claude-sdk-adapter/tests/native_model_options.test.mjs @@ -2,7 +2,7 @@ import assert from "node:assert/strict"; import test from "node:test"; import { parseRequest } from "../dist/request.js"; -const request = native_model_options => JSON.stringify({ type: "executor_prepare", model: "fixture", system_prompt: "", cwd: process.cwd(), native_model_options }); +const request = native_model_options => JSON.stringify({ type: "executor_prepare", preparation_deadline:Date.now()+60000, model: "fixture", system_prompt: "", cwd: process.cwd(), native_model_options }); test("private bridge accepts compiled native options and rejects malformed structure", () => { const native = { effort: "high", thinking: { type: "enabled", budgetTokens: 1024, display: "omitted" } }; assert.deepEqual(parseRequest(request(native)).native_model_options, native); @@ -12,5 +12,12 @@ test("private bridge accepts compiled native options and rejects malformed struc }); test("private bridge rejects the public configuration field", () => { - assert.throws(() => parseRequest(JSON.stringify({type:"executor_prepare", model:"fixture",system_prompt:"",cwd:process.cwd(),harness_config:{}})), {message:"invalid_request"}); + assert.throws(() => parseRequest(JSON.stringify({type:"executor_prepare",preparation_deadline:Date.now()+60000, model:"fixture",system_prompt:"",cwd:process.cwd(),harness_config:{}})), {message:"invalid_request"}); +}); + +test("Executor preparation requires an absolute integer deadline and other requests cannot carry one", () => { + const base={type:"executor_prepare",model:"fixture",system_prompt:"",cwd:process.cwd()}; + for(const value of [undefined,null,0,-1,1.5,"1",Number.MAX_SAFE_INTEGER+1])assert.throws(()=>parseRequest(JSON.stringify({...base,preparation_deadline:value})),/invalid_request/); + assert.equal(parseRequest(JSON.stringify({...base,preparation_deadline:1})).preparation_deadline,1); + assert.throws(()=>parseRequest(JSON.stringify({...base,type:"start",input:[{content:[{type:"input_text",text:"x"}]}],preparation_deadline:1})),/invalid_request/); }); diff --git a/packages/claude-sdk-adapter/tests/startup_budget.test.mjs b/packages/claude-sdk-adapter/tests/startup_budget.test.mjs index a1a555ce8..77b9c7f31 100644 --- a/packages/claude-sdk-adapter/tests/startup_budget.test.mjs +++ b/packages/claude-sdk-adapter/tests/startup_budget.test.mjs @@ -20,7 +20,7 @@ lines.on("line",line=>{ timer=setTimeout(()=>{ process.send({kind:"initialized",elapsed:performance.now()-started}); reply(message.request_id,{commands:[],models:[],account:{},hooks_applied:true}); - },16000); + },Number(process.argv[1])); } else if(message.type==="control_request" && message.request.subtype==="mcp_status") { process.send({kind:"required_status"}); reply(message.request_id,{mcpServers:[{name:"fixture",status:"connected",tools:[{name:"echo"}]}]}); @@ -33,7 +33,7 @@ lines.on("close",()=>{clearTimeout(timer);process.disconnect();}); const childProcess = ` import { spawn as nativeSpawn } from "node:child_process"; export function spawn(command,args,options) { - const child=nativeSpawn(process.execPath,["--input-type=module","-e",${JSON.stringify(native)}],{...options,stdio:[...options.stdio,"ipc"]}); + const child=nativeSpawn(process.execPath,["--input-type=module","-e",${JSON.stringify(native)},process.argv[1]],{...options,stdio:[...options.stdio,"ipc"]}); process.send({kind:"native_spawned",pid:child.pid}); child.on("message",message=>process.send(message)); child.once("close",()=>process.send({kind:"native_closed"})); @@ -42,7 +42,11 @@ export function spawn(command,args,options) { `; const bridge = ` import { registerHooks } from "node:module"; +import { startup } from ${JSON.stringify(new URL("../node_modules/@anthropic-ai/claude-agent-sdk/sdk.mjs", import.meta.url).href)}; +globalThis.actualStartup=args=>{process.send({kind:"budget",remaining:args.initializeTimeoutMs,now:Date.now()});return startup(args)}; registerHooks({resolve(specifier,context,next){ + if(specifier==="@anthropic-ai/claude-agent-sdk" && context.parentURL===${JSON.stringify(new URL("../dist/adapter.js", import.meta.url).href)}) + return {url:"data:text/javascript,"+encodeURIComponent("export { getSessionInfo,query } from " + ${JSON.stringify(JSON.stringify(new URL("../node_modules/@anthropic-ai/claude-agent-sdk/sdk.mjs", import.meta.url).href))} + ";export const startup=globalThis.actualStartup;"),shortCircuit:true}; if(specifier==="node:child_process" && context.parentURL===${JSON.stringify(new URL("../dist/native.js", import.meta.url).href)}) return {url:"data:text/javascript,"+encodeURIComponent(${JSON.stringify(childProcess)}),shortCircuit:true}; return next(specifier,context); @@ -51,15 +55,16 @@ await import(${JSON.stringify(new URL("../dist/main.js", import.meta.url).href)} process.disconnect(); `; -test("the SDK accepts delayed initialization on both adapter startup paths", {timeout:45000,concurrency:true}, async t => { - await Promise.all(["executor", "required-mcp"].map(mode => t.test(mode, async t => { +test("the SDK accepts delayed initialization on both adapter startup paths", {timeout:100000,concurrency:true}, async t => { + await Promise.all(["executor", "required-mcp", "short", "expired", "cancel"].map(mode => t.test(mode, async t => { const parent=join(homedir(),".oac","tests"); mkdirSync(parent,{recursive:true}); const root=mkdtempSync(join(parent,"claude-startup-budget-")); - const child=spawn(process.execPath,["--input-type=module","-e",bridge],{ + const child=spawn(process.execPath,["--input-type=module","-e",bridge,mode==="required-mcp"?"16000":"61000"],{ cwd:root,env:{PATH:process.env.PATH,HOME:root,CLAUDE_CONFIG_DIR:root},stdio:["pipe","pipe","pipe","ipc"], }); const events=[],observations=[]; + const deadline=Date.now()+(mode==="expired"?-1:mode==="short"?2000:90000); let stderr="",nativePID; child.stderr.on("data",value=>{stderr+=value}); const closed=new Promise(resolve=>child.once("close",(code,signal)=>resolve({code,signal}))); @@ -73,26 +78,36 @@ test("the SDK accepts delayed initialization on both adapter startup paths", {ti }; createInterface({input:child.stdout}).on("line",line=>{ const event=JSON.parse(line);events.push(event); - if(event.type==="error")reject(new Error("Preparation failed: "+event.code)); + if(event.type==="error") { if(["short","expired","cancel"].includes(mode))resolve();else reject(new Error("Preparation failed: "+event.code)); } notify(); }); child.on("message",message=>{ observations.push(message); - if(message.kind==="native_spawned")nativePID=message.pid; + if(message.kind==="native_spawned"){nativePID=message.pid;if(mode==="cancel")child.kill("SIGTERM");} if(message.kind==="native_closed")nativePID=undefined; notify(); }); - child.once("close",()=>reject(new Error("Bridge closed before readiness: "+stderr))); + child.once("close",()=>mode==="cancel"?resolve():reject(new Error("Bridge closed before readiness: "+stderr))); }); - child.stdin.write(JSON.stringify({type:mode==="executor"?"executor_prepare":"start",cwd:root,model:"fixture",system_prompt:"", + child.stdin.write(JSON.stringify({type:mode==="required-mcp"?"start":"executor_prepare",cwd:root,model:"fixture",system_prompt:"", + ...(mode!=="required-mcp"?{preparation_deadline:deadline}:{}), ...(mode==="required-mcp"?{input:[{content:[{type:"input_text",text:"held until initialization"}]}], mcp_http_servers:[{server_label:"fixture",server_url:"https://example.invalid/mcp",allowed_tools:["echo"],required:true}]}:{})})+"\n"); await ready; - assert.ok(observations.find(event=>event.kind==="initialized").elapsed>=15000); - assert.deepEqual(observations.map(event=>event.kind),mode==="executor"?["native_spawned","initialized"]:["native_spawned","initialized","required_status","input"]); + if (["short","expired","cancel"].includes(mode)) { + assert.equal(events.some(event=>event.type==="executor_ready"),false); + assert.equal(observations.some(event=>event.kind==="input"),false); + assert.equal(observations.some(event=>event.kind==="initialized"),false); + if(mode==="cancel")assert.deepEqual(events,[]);else assert.equal(events.at(-1).code,"execution_failed"); + } else assert.ok(observations.find(event=>event.kind==="initialized").elapsed >= (mode==="executor"?60000:15000)); + const budget=observations.find(event=>event.kind==="budget"); + if(mode==="expired")assert.equal(budget,undefined); + else if(mode!=="required-mcp") {assert.ok(budget.remaining>0);assert.ok(Math.abs(budget.now+budget.remaining-deadline)<20);} + + if (!["short","expired","cancel"].includes(mode)) assert.deepEqual(observations.filter(event=>event.kind!=="budget").map(event=>event.kind),mode==="executor"?["native_spawned","initialized"]:["native_spawned","initialized","required_status","input"]); if(mode==="executor")assert.deepEqual(events,[{type:"executor_ready",protocol:3}]); child.stdin.end(); assert.deepEqual(await closed,{code:0,signal:null},stderr); - assert.equal(observations.filter(event=>event.kind==="native_closed").length,1); + assert.equal(observations.filter(event=>event.kind==="native_closed").length,mode==="expired"?0:1); }))); });