Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
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
29 changes: 16 additions & 13 deletions apps/daemon/internal/dispatch/runtime_preparation.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,19 +13,18 @@ import (
"github.com/google/uuid"
)

const runtimePreparationTimeout = 120 * time.Second

// Router.mu protects one connection-local transfer. Partial installation data
// belongs to the bound Environment and is never removed by transfer cleanup.
type runtimePreparationTransfer struct {
envelope proto.Envelope
request proto.RuntimePreparePayload
data []byte
ready chan struct{}
cancel context.CancelFunc
finished bool
apply bool
uncertain bool
transferDeadline time.Time
envelope proto.Envelope
request proto.RuntimePreparePayload
data []byte
ready chan struct{}
cancel context.CancelFunc
finished bool
apply bool
uncertain bool
}

func (r *Router) handleRuntimePrepare(ctx context.Context, env proto.Envelope) error {
Expand Down Expand Up @@ -65,9 +64,10 @@ func (r *Router) handleRuntimePrepare(ctx context.Context, env proto.Envelope) e
r.mu.Unlock()
return r.sendRuntimePrepareResult(ctx, env.ID, rejectedRuntimePreparation("resource_unavailable"))
}
owner, cancel := context.WithTimeout(context.WithoutCancel(ctx), runtimePreparationTimeout)
owner, cancel := context.WithTimeout(context.WithoutCancel(ctx), time.Duration(request.BudgetMS)*time.Millisecond)
u := &runtimePreparationTransfer{
envelope: env, request: request, data: make([]byte, 0, request.SizeBytes),
transferDeadline: time.Now().Add(time.Duration(proto.RuntimePrepareTransferBudgetMS) * time.Millisecond),
envelope: env, request: request, data: make([]byte, 0, request.SizeBytes),
ready: make(chan struct{}), cancel: cancel,
}
r.runtimePreparation = u
Expand Down Expand Up @@ -100,7 +100,7 @@ func (r *Router) handleRuntimePrepare(ctx context.Context, env proto.Envelope) e
return nil
}
apply := false
if request.Step == "commit" && len(u.data) == u.request.SizeBytes {
if request.Step == "commit" && len(u.data) == u.request.SizeBytes && time.Now().Before(u.transferDeadline) {
if u.request.Action == "finalize" || u.request.Action == "initialize" {
apply = true
} else {
Expand Down Expand Up @@ -135,7 +135,10 @@ func (r *Router) finishRuntimePreparationTransferLocked(u *runtimePreparationTra
func (r *Router) runRuntimePreparationTransfer(ctx context.Context, u *runtimePreparationTransfer, apply func(context.Context, proto.RuntimePreparePayload, []byte) error) {
defer r.shutdownWG.Done()
defer u.cancel()
staging := time.NewTimer(time.Until(u.transferDeadline))
defer staging.Stop()
select {
case <-staging.C:
case <-u.ready:
case <-r.shutdownCh:
case <-ctx.Done():
Expand Down
81 changes: 79 additions & 2 deletions apps/daemon/internal/dispatch/runtime_preparation_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import (
"os"
"path/filepath"
"testing"
"testing/synctest"
"time"

"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent"
Expand Down Expand Up @@ -57,7 +58,7 @@ func capabilityEnvelope(t *testing.T, id string, request proto.RuntimePreparePay
func capabilityBegin(environment, session string, body []byte) proto.RuntimePreparePayload {
digest := sha256.Sum256(body)
return proto.RuntimePreparePayload{
Step: "begin", Action: "skill", EnvironmentID: environment, SessionID: session,
BudgetMS: 300000, Step: "begin", Action: "skill", EnvironmentID: environment, SessionID: session,
Skill: &agentskill.Metadata{Type: "inline", Name: "proof", Description: "A proof."},
SizeBytes: len(body), SHA256: hex.EncodeToString(digest[:]),
}
Expand Down Expand Up @@ -247,7 +248,7 @@ func TestRuntimePreparationCancellationKeepsOwnershipUntilApplyStops(t *testing.
r, sender, environment, session := capabilitiesTestRouter(t)
ctx, cancel := context.WithCancel(context.Background())
id := uuid.NewString()
request := proto.RuntimePreparePayload{Step: "begin", Action: "finalize", EnvironmentID: environment, SessionID: session, Sources: &agentcapabilities.Input{}}
request := proto.RuntimePreparePayload{BudgetMS: 300000, Step: "begin", Action: "finalize", EnvironmentID: environment, SessionID: session, Sources: &agentcapabilities.Input{}}
owner := &runtimePreparationTransfer{envelope: capabilityEnvelope(t, id, request), request: request, ready: make(chan struct{}), cancel: cancel, finished: true, apply: true}
close(owner.ready)
r.runtimePreparation = owner
Expand Down Expand Up @@ -346,3 +347,79 @@ func TestRuntimePreparationResultCategoriesAndUnknownOwnership(t *testing.T) {
t.Fatal("shutdown claimed uncertain mutation settled")
}
}

func TestRuntimePreparationStagingDeadlineRejectsWithoutApplying(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
r, sender, environment, session := capabilitiesTestRouter(t)
id := uuid.NewString()
request := capabilityBegin(environment, session, []byte("abc"))
if err := r.Handle(t.Context(), capabilityEnvelope(t, id, request)); err != nil {
t.Fatal(err)
}
capabilitiesReceipt(t, sender, id, "ready")
time.Sleep(121 * time.Second)
synctest.Wait()
capabilitiesReceipt(t, sender, id, "rejected")
r.mu.Lock()
pending := r.runtimePreparation != nil
r.mu.Unlock()
if pending {
t.Fatal("uncommitted staging retained ownership")
}
shutdownCapabilitiesRouter(t, r)
})
}

func TestRuntimePreparationApplyOutlivesStagingAndHonorsBudget(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
r, sender, environment, session := capabilitiesTestRouter(t)
ctx, cancel := context.WithTimeout(t.Context(), 5*time.Minute)
id := uuid.NewString()
request := capabilityBegin(environment, session, []byte("abc"))
owner := &runtimePreparationTransfer{envelope: capabilityEnvelope(t, id, request), request: request, ready: make(chan struct{}), cancel: cancel, finished: true, apply: true, transferDeadline: time.Now().Add(2 * time.Minute)}
close(owner.ready)
r.runtimePreparation = owner
r.shutdownWG.Add(1)
stopped := make(chan struct{})
go r.runRuntimePreparationTransfer(ctx, owner, func(ctx context.Context, _ proto.RuntimePreparePayload, _ []byte) error {
<-ctx.Done()
close(stopped)
return ctx.Err()
})
time.Sleep(121 * time.Second)
synctest.Wait()
select {
case <-stopped:
t.Fatal("staging timeout killed apply")
default:
}
time.Sleep(180 * time.Second)
synctest.Wait()
<-stopped
capabilitiesReceipt(t, sender, id, "unknown")
// Unknown effects remain retained; this test must not claim safe replay.
r.mu.Lock()
unknown := r.runtimePreparation == owner && owner.uncertain
r.mu.Unlock()
if !unknown {
t.Fatal("unknown apply lost ownership")
}
})
}

func TestRuntimePreparationHonorsBeginBudgetDuringStaging(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
r, sender, environment, session := capabilitiesTestRouter(t)
id := uuid.NewString()
request := capabilityBegin(environment, session, []byte("abc"))
request.BudgetMS = 10
if err := r.Handle(t.Context(), capabilityEnvelope(t, id, request)); err != nil {
t.Fatal(err)
}
capabilitiesReceipt(t, sender, id, "ready")
time.Sleep(11 * time.Millisecond)
synctest.Wait()
capabilitiesReceipt(t, sender, id, "rejected")
shutdownCapabilitiesRouter(t, r)
})
}
Original file line number Diff line number Diff line change
Expand Up @@ -101,15 +101,13 @@ func initializationEnvironment(configured map[string]string) []string {
// The shared process owner settles the leader and descendants. Readers finish
// before return; output is discarded with constant memory, never put in errors.
func runInitializationProcess(ctx context.Context, binary string, args []string, directory string, env []string) error {
operation, cancel := context.WithTimeout(ctx, 30*time.Minute)
defer cancel()
if operation.Err() != nil {
if ctx.Err() != nil {
return ErrInitializationUnconfirmed
}
if info, err := os.Stat(directory); err != nil || !info.IsDir() {
return &InitializationFailure{}
}
process, err := clirunner.Start(clirunner.StartOptions{Parent: operation, Binary: binary, Args: args,
process, err := clirunner.Start(clirunner.StartOptions{Parent: ctx, Binary: binary, Args: args,
Dir: directory, Env: env, KillTimeout: 250 * time.Millisecond})
if err != nil {
return ErrInitializationUnconfirmed
Expand All @@ -127,7 +125,7 @@ func runInitializationProcess(ctx context.Context, binary string, args []string,
}
first, second := <-finished, <-finished
_ = process.Wait()
if operation.Err() != nil || first != nil || second != nil || process.Cmd.ProcessState == nil {
if ctx.Err() != nil || first != nil || second != nil || process.Cmd.ProcessState == nil {
return ErrInitializationUnconfirmed
}
code := process.Cmd.ProcessState.ExitCode()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -348,7 +348,7 @@ func TestRuntimePreparationRejectsFilesAfterFinalization(t *testing.T) {
t.Fatal(err)
}
digest := sha256.Sum256(nil)
input := proto.RuntimePreparePayload{Step: "begin", EnvironmentID: b.environment, SessionID: b.capabilityIdentity().SessionID,
input := proto.RuntimePreparePayload{BudgetMS: 300000, Step: "begin", EnvironmentID: b.environment, SessionID: b.capabilityIdentity().SessionID,
Action: "file", File: &proto.RuntimeInitialFile{Path: "/workspace/file"}, SHA256: hex.EncodeToString(digest[:])}
if err = b.ApplyRuntimePreparation(t.Context(), input, nil); !errors.Is(err, agentcapabilities.ErrInvalid) {
t.Fatal("finalized Runtime accepted file", err)
Expand All @@ -361,3 +361,25 @@ func TestRuntimePreparationRejectsFilesAfterFinalization(t *testing.T) {
t.Fatal("finalized Runtime accepted initialize", err)
}
}

func TestRuntimeInitializationDeadlineStopsDescendants(t *testing.T) {
binary, err := os.Executable()
if err != nil {
t.Fatal(err)
}
directory := t.TempDir()
ctx, cancel := context.WithTimeout(t.Context(), 500*time.Millisecond)
defer cancel()
env := append(initializationEnvironment(nil), "OAC_INITIALIZATION_FIXTURE=parent", "OAC_INITIALIZATION_FIXTURE_DIR="+directory)
err = runInitializationProcess(ctx, binary, []string{"-test.run=^TestRuntimeInitializationChild$"}, directory, env)
if !errors.Is(err, ErrInitializationUnconfirmed) || !errors.Is(ctx.Err(), context.DeadlineExceeded) {
t.Fatal("deadline outcome", err, ctx.Err())
}
if _, err := os.Stat(filepath.Join(directory, "started")); err != nil {
t.Fatal("fixture did not start", err)
}
time.Sleep(1100 * time.Millisecond)
if _, err := os.Stat(filepath.Join(directory, "late")); !os.IsNotExist(err) {
t.Fatal("descendant survived deadline")
}
}
4 changes: 2 additions & 2 deletions apps/daemon/internal/localworkspace/snapshot_marker_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -186,7 +186,7 @@ func TestSnapshotMarkerRefusesDamagedState(t *testing.T) {
func TestCapabilityFinalizeRecordsCompletionAndRejectsLaterImports(t *testing.T) {
b, _ := markerBinding(t)
identity := b.capabilityIdentity()
finalize := proto.RuntimePreparePayload{Step: "begin", EnvironmentID: identity.EnvironmentID, SessionID: identity.SessionID, Action: "finalize", Sources: &agentcapabilities.Input{}}
finalize := proto.RuntimePreparePayload{BudgetMS: 300000, Step: "begin", EnvironmentID: identity.EnvironmentID, SessionID: identity.SessionID, Action: "finalize", Sources: &agentcapabilities.Input{}}
if err := b.ApplyRuntimePreparation(t.Context(), finalize, nil); err != nil {
t.Fatal(err)
}
Expand All @@ -197,7 +197,7 @@ func TestCapabilityFinalizeRecordsCompletionAndRejectsLaterImports(t *testing.T)
if err := os.RemoveAll(b.capabilityRoot); err != nil {
t.Fatal(err)
}
skill := proto.RuntimePreparePayload{Step: "begin", EnvironmentID: identity.EnvironmentID, SessionID: identity.SessionID, Action: "skill", Skill: &agentskill.Metadata{Type: "inline", Name: "example", Description: "Example"}, SizeBytes: 1, SHA256: "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"}
skill := proto.RuntimePreparePayload{BudgetMS: 300000, Step: "begin", EnvironmentID: identity.EnvironmentID, SessionID: identity.SessionID, Action: "skill", Skill: &agentskill.Metadata{Type: "inline", Name: "example", Description: "Example"}, SizeBytes: 1, SHA256: "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"}
if err := b.ApplyRuntimePreparation(t.Context(), skill, []byte("x")); err == nil {
t.Fatal("completed identity accepted import")
}
Expand Down
2 changes: 1 addition & 1 deletion contracts/agents-api/environments.md
Original file line number Diff line number Diff line change
Expand Up @@ -229,7 +229,7 @@ An Environment's initialization is `pending`, `running`, `complete` or `failed`,

- The Worker's initialization scheduler scans 32 Environments at a time, wraps at the end and bounds concurrent preparations by execution concurrency, independently of Provider maintenance.
- A missing socket does not consume a pending attempt. An unavailable Harness fails before installation. Each operation rechecks current authority and the original socket; completion rechecks the exact binding.
- Each file transfer, configure, Skill, Plugin, package and setup step has a two-minute budget; the whole initialization has 30 minutes. Initial input keeps its five-minute admission deadline, so large installations should start from an idle Session.
- Each file, Skill or Plugin transfer has a two-minute staging budget. Configure, package installation, setup commands and snapshot finalization share the whole initialization’s 30-minute budget, whose remaining duration Core carries on each operation and the Runtime enforces. Initial input keeps its five-minute admission deadline, so large installations should start from an idle Session.
- A running initialization whose owner is lost, including across a Core restart, fails as unconfirmed; nothing is replayed. A completed Environment never reinstalls on reconnect or native recovery, so later user changes survive.
- Failure is terminal for the Session but destroys neither compute nor files.

Expand Down
4 changes: 2 additions & 2 deletions contracts/agents-api/zh/environments.md
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
---
title: "环境与模板"
source: contracts/agents-api/environments.md
source_hash: 48b9d8679a983adc5eb9164c9a94e1137bf75800fbb1a0f81a0ae1b3cfccdafb
source_hash: 9007be1099cf50f5ba376d2de4a060d039b6ee0b1c13c78a94665c28c2fc3fc6
---

Environment 是 Session 的执行资源,包括 Harness 运行所在的机器、工作区以及已完成准备的能力。Session 通过其 `environment` 配置创建 Environment;不存在独立的 create 调用。Environment Template 是 Session 创建时解析的可复用准备配置。本契约涵盖这两类资源、两种放置方式、输入接纳、能力准备、Skills、Plugins 和 MCP 连接来源。
Expand Down Expand Up @@ -231,7 +231,7 @@ Environment 的初始化状态为 `pending`、`running`、`complete` 或 `failed

- Worker 的初始化调度器每次扫描 32 个 Environment,在末尾循环回绕,并依据执行并发度限制并发准备,且独立于 Provider 维护。
- 缺少套接字不会消耗一次 pending 尝试。Harness 不可用时,会在安装前失败。每个操作都会重新检查当前权限和原始套接字;完成时还会重新检查精确绑定。
- 每个文件传输、configure、Skill、Plugin、软件包和设置步骤都有两分钟的预算;整个初始化过程有 30 分钟。初始输入仍保留其五分钟接纳期限,因此大型安装应从空闲 Session 开始。
- 每个文件、Skill 或 Plugin 传输的暂存预算为两分钟。configure、软件包安装、设置命令和快照最终确定共享整个初始化过程的 30 分钟预算,Core 在每项操作中携带剩余时长,由 Runtime 执行。初始输入仍保留其五分钟接纳期限,因此大型安装应从空闲 Session 开始。
- 正在运行且所有权丧失的初始化,包括跨 Core 重启丧失所有权,会作为未确认而失败;不会重播任何内容。已完成的 Environment 在重连或原生恢复时绝不会重新安装,因此用户后续更改会保留下来。
- 失败对 Session 而言是终止状态,但不会销毁计算资源或文件。

Expand Down
2 changes: 1 addition & 1 deletion docs/runtime-protocol.md
Original file line number Diff line number Diff line change
Expand Up @@ -110,7 +110,7 @@ Usage frames and the final usage snapshot each carry the cumulative measurement

## Preparation and execution order

Environment initialization uses `runtime_prepare` on every connection, managed or user-owned; the [Environment contract](../contracts/agents-api/environments.md#runtime-capability-preparation) owns what is prepared and when. For a file or archive, send `begin`, wait for `ready`, send ordered chunks and await each matching `received` offset, then send `commit` and await `completed`. Initialization and finalization have typed headers without file data. Validate the expected outcome, offset, size and finite error code with the shared validator. One transfer is allowed per connection. A chunk receipt confirms staged bytes, not installation; a completed commit confirms that operation, not that a later Turn ran.
Environment initialization uses `runtime_prepare` on every connection, managed or user-owned; the [Environment contract](../contracts/agents-api/environments.md#runtime-capability-preparation) owns what is prepared and when. For a file or archive, send `begin`, wait for `ready`, send ordered chunks and await each matching `received` offset, then send `commit` and await `completed`. Initialization and finalization have typed headers without file data. Every `begin` requires `budget_ms`, an integer from 1 through 1800000 carrying the remaining Core-owned initialization budget; chunks and commits omit it. The Runtime enforces that relative budget from receipt of `begin`, and Core waits no longer than its original operation deadline. The complete `begin`–`commit` staging phase has a separate two-minute limit; applying a committed operation and waiting for `completed` use the remaining initialization budget, not the staging limit. Expiry before application rejects the transfer; expiry during application reports `unknown` after local mutations stop and never authorizes replay. Validate the expected outcome, offset, size and finite error code with the shared validator. One transfer is allowed per connection. A chunk receipt confirms staged bytes, not installation; a completed commit confirms that operation, not that a later Turn ran.

An execution Turn runs in five steps:

Expand Down
Loading
Loading