From 5718fc11e60bcfdf6ed43de29742700d38024e02 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Sat, 10 Oct 2026 16:37:30 +0000 Subject: [PATCH] Honor initialization execution budget --- .../internal/dispatch/runtime_preparation.go | 29 ++++--- .../dispatch/runtime_preparation_test.go | 81 ++++++++++++++++++- .../runtime_initialization_process.go | 8 +- .../runtime_initialization_test.go | 24 +++++- .../localworkspace/snapshot_marker_test.go | 4 +- contracts/agents-api/environments.md | 2 +- contracts/agents-api/zh/environments.md | 4 +- docs/runtime-protocol.md | 2 +- docs/zh/runtime-protocol.md | 4 +- internal/agentdaemon/proto/runtime_prepare.go | 18 +++-- .../agentdaemon/proto/runtime_prepare_test.go | 26 +++++- internal/agentdaemon/proto/version.go | 2 +- .../execution/runtime_capabilities_test.go | 6 +- .../execution/runtime_initialization.go | 15 ++-- .../core/internal/execution/runtime_setup.go | 7 ++ .../internal/execution/runtime_setup_test.go | 40 ++++++++- .../runtimegateway/runtime_prepare.go | 12 ++- .../runtimegateway/runtime_prepare_test.go | 78 +++++++++++++++++- 18 files changed, 302 insertions(+), 60 deletions(-) diff --git a/apps/daemon/internal/dispatch/runtime_preparation.go b/apps/daemon/internal/dispatch/runtime_preparation.go index 1c777f0fe..926ddc5f9 100644 --- a/apps/daemon/internal/dispatch/runtime_preparation.go +++ b/apps/daemon/internal/dispatch/runtime_preparation.go @@ -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 { @@ -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 @@ -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 { @@ -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(): diff --git a/apps/daemon/internal/dispatch/runtime_preparation_test.go b/apps/daemon/internal/dispatch/runtime_preparation_test.go index 7a8a31c70..90a0c418e 100644 --- a/apps/daemon/internal/dispatch/runtime_preparation_test.go +++ b/apps/daemon/internal/dispatch/runtime_preparation_test.go @@ -9,6 +9,7 @@ import ( "os" "path/filepath" "testing" + "testing/synctest" "time" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" @@ -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[:]), } @@ -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 @@ -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) + }) +} diff --git a/apps/daemon/internal/localworkspace/runtime_initialization_process.go b/apps/daemon/internal/localworkspace/runtime_initialization_process.go index 68174485d..a8464ed97 100644 --- a/apps/daemon/internal/localworkspace/runtime_initialization_process.go +++ b/apps/daemon/internal/localworkspace/runtime_initialization_process.go @@ -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 @@ -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() diff --git a/apps/daemon/internal/localworkspace/runtime_initialization_test.go b/apps/daemon/internal/localworkspace/runtime_initialization_test.go index e16e46787..47896582a 100644 --- a/apps/daemon/internal/localworkspace/runtime_initialization_test.go +++ b/apps/daemon/internal/localworkspace/runtime_initialization_test.go @@ -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) @@ -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") + } +} diff --git a/apps/daemon/internal/localworkspace/snapshot_marker_test.go b/apps/daemon/internal/localworkspace/snapshot_marker_test.go index 7f610e721..90132e1e5 100644 --- a/apps/daemon/internal/localworkspace/snapshot_marker_test.go +++ b/apps/daemon/internal/localworkspace/snapshot_marker_test.go @@ -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) } @@ -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") } diff --git a/contracts/agents-api/environments.md b/contracts/agents-api/environments.md index 099243c5f..79126cbc7 100644 --- a/contracts/agents-api/environments.md +++ b/contracts/agents-api/environments.md @@ -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. diff --git a/contracts/agents-api/zh/environments.md b/contracts/agents-api/zh/environments.md index 570681196..ffed2f938 100644 --- a/contracts/agents-api/zh/environments.md +++ b/contracts/agents-api/zh/environments.md @@ -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 连接来源。 @@ -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 而言是终止状态,但不会销毁计算资源或文件。 diff --git a/docs/runtime-protocol.md b/docs/runtime-protocol.md index b99b6b668..e7c818fa8 100644 --- a/docs/runtime-protocol.md +++ b/docs/runtime-protocol.md @@ -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: diff --git a/docs/zh/runtime-protocol.md b/docs/zh/runtime-protocol.md index 92de66cca..ab1b54296 100644 --- a/docs/zh/runtime-protocol.md +++ b/docs/zh/runtime-protocol.md @@ -1,7 +1,7 @@ --- title: "Core–Runtime 协议" source: docs/runtime-protocol.md -source_hash: 59633fbfa54d9071ed6a5ef510328d9216f4d8572a0cc068de65c93ecabefc59 +source_hash: 142a3d2a3dc5b790bcc80dff2364d094857f835e0e0687f91373251a4f4d9b3b --- 此协议在 Runtime daemon 获取机器凭据后连接 Core 与 daemon,定义 daemon 连接上消息的含义和顺序。wire 类型、限制和验证器仅在 [`internal/agentdaemon/proto`](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/internal/agentdaemon/proto) 中定义一次;Core 的 [gateway](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/services/core/internal/runtimegateway) 与参考 Runtime 的 [dispatcher](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/apps/daemon/internal/dispatch) 都使用它们,因此无需同步第二套 payload schema。签发凭据和打开连接的 HTTP 路由见[机器连接 API](../../contracts/agents-api/zh/machine-api.md)。 @@ -112,7 +112,7 @@ Usage frame 和最终 usage snapshot 都携带当前执行的累计测量,替 ## 准备与执行顺序 {#preparation-and-execution-order} -无论托管还是用户自有环境,每条连接都通过 `runtime_prepare` 初始化 Environment;[Environment 契约](../../contracts/agents-api/zh/environments.md#runtime-capability-preparation)负责准备内容和时机。传输文件或 archive 时,发送 `begin`,等待 `ready`,发送有序 chunk 并等待每个匹配的 `received` offset,再发送 `commit` 并等待 `completed`。初始化和终结阶段使用不含文件数据的类型化 header。使用共享 validator 验证预期结果、offset、size 和有限错误 code。每条连接允许一个 transfer。chunk 回执确认暂存字节,不确认安装;完成的 commit 确认该操作,不证明后续 Turn 已运行。 +无论托管还是用户自有环境,每条连接都通过 `runtime_prepare` 初始化 Environment;[Environment 契约](../../contracts/agents-api/zh/environments.md#runtime-capability-preparation)负责准备内容和时机。传输文件或 archive 时,发送 `begin`,等待 `ready`,发送有序 chunk 并等待每个匹配的 `received` offset,再发送 `commit` 并等待 `completed`。初始化和终结阶段使用不含文件数据的类型化 header。每个 `begin` 必须携带 `budget_ms`,取值为 1 至 1800000 的整数,表示 Core 所拥有的初始化预算剩余毫秒数;chunk 和 commit 不携带该字段。Runtime 从收到 `begin` 起执行此相对预算,Core 的等待不超过原操作期限。完整的 `begin`–`commit` 暂存阶段另有两分钟限制;应用已提交操作及等待 `completed` 使用剩余初始化预算,不受暂存期限限制。应用前到期会拒绝传输;应用中到期则在本地变更停止后报告 `unknown`,绝不授权重放。使用共享 validator 验证预期结果、offset、size 和有限错误 code。每条连接允许一个 transfer。chunk 回执确认暂存字节,不确认安装;完成的 commit 确认该操作,不证明后续 Turn 已运行。 执行 Turn 分为五步: diff --git a/internal/agentdaemon/proto/runtime_prepare.go b/internal/agentdaemon/proto/runtime_prepare.go index bbb2549c6..576cd8be4 100644 --- a/internal/agentdaemon/proto/runtime_prepare.go +++ b/internal/agentdaemon/proto/runtime_prepare.go @@ -14,11 +14,13 @@ import ( ) const ( - TypeRuntimePrepare = "runtime_prepare" - TypeRuntimePrepareResult = "runtime_prepare_result" - RuntimePrepareMaxBytes = 50 << 20 - RuntimePrepareChunkBytes = WorkspaceWriteChunkBytes - RuntimePrepareMaxFrameBytes = 1 << 20 + TypeRuntimePrepare = "runtime_prepare" + TypeRuntimePrepareResult = "runtime_prepare_result" + RuntimePrepareMaxBudgetMS = 30 * 60 * 1000 + RuntimePrepareTransferBudgetMS = 2 * 60 * 1000 + RuntimePrepareMaxBytes = 50 << 20 + RuntimePrepareChunkBytes = WorkspaceWriteChunkBytes + RuntimePrepareMaxFrameBytes = 1 << 20 ) // RuntimeInitialFile addresses a file within the logical workspace. @@ -38,6 +40,7 @@ type RuntimeInitialization struct { // Runtime resolves logical paths and owns installation destinations. Envelope.ID // identifies one connection-local transfer. type RuntimePreparePayload struct { + BudgetMS int64 `json:"budget_ms,omitempty"` Step string `json:"step"` EnvironmentID string `json:"environment_id,omitempty"` SessionID string `json:"session_id,omitempty"` @@ -64,6 +67,9 @@ type RuntimePrepareResultPayload struct { func ValidRuntimePrepareRequest(p RuntimePreparePayload) bool { if p.Step == "begin" { + if p.BudgetMS <= 0 || p.BudgetMS > RuntimePrepareMaxBudgetMS { + return false + } for _, id := range []string{p.EnvironmentID, p.SessionID} { value, err := uuid.Parse(id) if err != nil || value == uuid.Nil || value.String() != id { @@ -119,7 +125,7 @@ func ValidRuntimePrepareRequest(p RuntimePreparePayload) bool { encoded, err := json.Marshal(p) return err == nil && len(encoded) <= RuntimePrepareMaxFrameBytes } - if p.EnvironmentID != "" || p.SessionID != "" || p.Action != "" || p.Slot != 0 || + if p.BudgetMS != 0 || p.EnvironmentID != "" || p.SessionID != "" || p.Action != "" || p.Slot != 0 || p.Skill != nil || p.Plugin != nil || p.Sources != nil || p.File != nil || p.Initialization != nil || p.SizeBytes != 0 || p.SHA256 != "" { return false } diff --git a/internal/agentdaemon/proto/runtime_prepare_test.go b/internal/agentdaemon/proto/runtime_prepare_test.go index 84519b778..8b63322e9 100644 --- a/internal/agentdaemon/proto/runtime_prepare_test.go +++ b/internal/agentdaemon/proto/runtime_prepare_test.go @@ -13,7 +13,7 @@ import ( ) func capabilityBegin() RuntimePreparePayload { - return RuntimePreparePayload{Step: "begin", EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), + return RuntimePreparePayload{BudgetMS: 300000, Step: "begin", EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Action: "skill", Skill: &agentskill.Metadata{Type: "inline", Name: "example", Description: "Example"}, SizeBytes: 10, SHA256: strings.Repeat("a", 64)} } @@ -23,6 +23,9 @@ func TestCapabilitiesRequestValidation(t *testing.T) { t.Fatal("valid skill refused") } for name, mutate := range map[string]func(*RuntimePreparePayload){ + "missing budget": func(p *RuntimePreparePayload) { p.BudgetMS = 0 }, + "negative budget": func(p *RuntimePreparePayload) { p.BudgetMS = -1 }, + "excessive budget": func(p *RuntimePreparePayload) { p.BudgetMS = RuntimePrepareMaxBudgetMS + 1 }, "noncanonical identity": func(p *RuntimePreparePayload) { p.EnvironmentID = "AAAAAAAA-AAAA-4AAA-8AAA-AAAAAAAAAAAA" }, "zero identity": func(p *RuntimePreparePayload) { p.SessionID = uuid.Nil.String() }, "begin data": func(p *RuntimePreparePayload) { p.Data = []byte("x") }, @@ -91,7 +94,7 @@ func TestCapabilitiesRequestValidation(t *testing.T) { } func TestCapabilitiesFinalizeRequiresExplicitBoundedSelection(t *testing.T) { - p := RuntimePreparePayload{Step: "begin", EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Action: "finalize", Sources: &agentcapabilities.Input{}} + p := RuntimePreparePayload{BudgetMS: 300000, Step: "begin", EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Action: "finalize", Sources: &agentcapabilities.Input{}} if !ValidRuntimePrepareRequest(p) { t.Fatal("explicit empty finalization refused") } @@ -173,7 +176,7 @@ func TestLocalEnvironmentCarriesSelectionButNeverInstalledRoots(t *testing.T) { } func TestRuntimePreparationInitialActions(t *testing.T) { - base := RuntimePreparePayload{Step: "begin", EnvironmentID: uuid.NewString(), SessionID: uuid.NewString()} + base := RuntimePreparePayload{BudgetMS: 300000, Step: "begin", EnvironmentID: uuid.NewString(), SessionID: uuid.NewString()} for _, initialization := range []RuntimeInitialization{ {Action: "configure", Env: map[string]string{"EXAMPLE": "value"}}, {Action: "npm", Packages: []string{"typescript"}}, @@ -263,7 +266,7 @@ func TestRuntimePreparationExitCodes(t *testing.T) { } func TestRuntimePreparationPortableSources(t *testing.T) { - p := RuntimePreparePayload{Step: "begin", EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Action: "finalize", Sources: &agentcapabilities.Input{Directories: []string{`C:\Users\operator\skills`, `\\host\share\skills`}}} + p := RuntimePreparePayload{BudgetMS: 300000, Step: "begin", EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Action: "finalize", Sources: &agentcapabilities.Input{Directories: []string{`C:\Users\operator\skills`, `\\host\share\skills`}}} if !ValidRuntimePrepareRequest(p) { t.Fatal("portable sources rejected by wire") } @@ -271,3 +274,18 @@ func TestRuntimePreparationPortableSources(t *testing.T) { t.Fatal("Windows sources accepted by Linux resolver") } } + +func TestRuntimePreparationBudgetOnlyOnBegin(t *testing.T) { + for _, p := range []RuntimePreparePayload{{Step: "commit", BudgetMS: 1}, {Step: "chunk", Data: []byte("a"), BudgetMS: 1}} { + if ValidRuntimePrepareRequest(p) { + t.Fatal("budget admitted outside begin") + } + } + for _, budget := range []int64{1, RuntimePrepareMaxBudgetMS} { + p := capabilityBegin() + p.BudgetMS = budget + if !ValidRuntimePrepareRequest(p) { + t.Fatal("valid budget boundary rejected", budget) + } + } +} diff --git a/internal/agentdaemon/proto/version.go b/internal/agentdaemon/proto/version.go index 7f6f11783..114299f22 100644 --- a/internal/agentdaemon/proto/version.go +++ b/internal/agentdaemon/proto/version.go @@ -2,7 +2,7 @@ package proto // Version identifies the complete Core–Runtime wire contract. Change it when // removing or changing a payload or its semantics; deploy both endpoints together. -const Version = "0.14.0" +const Version = "0.15.0" // VersionCompatible accepts only this contract. Patch drift, prerelease suffixes // and malformed versions do not select an implicit compatibility path. diff --git a/services/core/internal/execution/runtime_capabilities_test.go b/services/core/internal/execution/runtime_capabilities_test.go index c2c29bca1..63319df84 100644 --- a/services/core/internal/execution/runtime_capabilities_test.go +++ b/services/core/internal/execution/runtime_capabilities_test.go @@ -40,7 +40,7 @@ func TestRuntimeCapabilitiesPreserveRawBundlesAndSetupOrdering(t *testing.T) { owner := agentcapabilities.Identity{EnvironmentID: uuid.NewString(), SessionID: uuid.NewString()} for _, op := range operations { peer := &capabilityFixture{outcome: "completed"} - if err := runRuntimeSetup(t.Context(), peer, owner, op); err != nil { + if err := runRuntimeSetup(initializationTestContext(t), peer, owner, op); err != nil { t.Fatal(err) } if _, err := uuid.Parse(peer.requestID); err != nil { @@ -64,14 +64,14 @@ func TestRuntimeCapabilitiesPreserveRawBundlesAndSetupOrdering(t *testing.T) { func TestRuntimeCapabilitiesConfirmedAndUnknownFailures(t *testing.T) { for _, outcome := range []string{"failed", "rejected", "unknown", "unexpected"} { peer := &capabilityFixture{outcome: outcome} - err := runRuntimeSetup(t.Context(), peer, agentcapabilities.Identity{}, runtimeSetupOperation{Request: proto.RuntimePreparePayload{Action: "finalize"}}) + err := runRuntimeSetup(initializationTestContext(t), peer, agentcapabilities.Identity{}, runtimeSetupOperation{Request: proto.RuntimePreparePayload{Action: "finalize"}}) var confirmed *runtimeStepFailure if err == nil || errors.As(err, &confirmed) != (outcome == "failed" || outcome == "rejected") { t.Fatal(outcome, err) } } peer := &capabilityFixture{outcome: "completed", err: errors.New(setupCanary)} - err := runRuntimeSetup(t.Context(), peer, agentcapabilities.Identity{}, runtimeSetupOperation{Request: proto.RuntimePreparePayload{}}) + err := runRuntimeSetup(initializationTestContext(t), peer, agentcapabilities.Identity{}, runtimeSetupOperation{Request: proto.RuntimePreparePayload{}}) if err == nil || bytes.Contains([]byte(err.Error()), []byte(setupCanary)) { t.Fatal("transport error leaked or succeeded", err) } diff --git a/services/core/internal/execution/runtime_initialization.go b/services/core/internal/execution/runtime_initialization.go index 8c2072578..1c958d8ea 100644 --- a/services/core/internal/execution/runtime_initialization.go +++ b/services/core/internal/execution/runtime_initialization.go @@ -8,6 +8,7 @@ import ( "time" "github.com/MiniMax-AI/OpenAgentCore/internal/agentcapabilities" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/environmentconfig" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" @@ -91,7 +92,7 @@ func (w *Worker) runEnvironmentInitializations(ctx context.Context) error { } func (w *Worker) initializeEnvironment(ctx context.Context, owner sessions.EnvironmentInitialization) { - operation, cancel := context.WithTimeout(ctx, 30*time.Minute) + operation, cancel := context.WithTimeout(ctx, time.Duration(proto.RuntimePrepareMaxBudgetMS)*time.Millisecond) defer cancel() failure := sessions.ProvisioningFailure{} err := w.prepareEnvironment(operation, owner, &failure) @@ -140,11 +141,10 @@ func (w *Worker) prepareEnvironment(ctx context.Context, owner sessions.Environm identity := agentcapabilities.Identity{EnvironmentID: owner.EnvironmentID, SessionID: owner.SessionID} operations := setupOperations(setup) for index := 0; index < len(cfg.Files)+len(operations); index++ { - step, stop := context.WithTimeout(ctx, 2*time.Minute) - err = w.lease.CheckOwnership(step) + err = w.lease.CheckOwnership(ctx) if err == nil { var currentPeer = peer - currentPeer, err = w.dispatcher.authorizedPeer(step, owner.DeviceID) + currentPeer, err = w.dispatcher.authorizedPeer(ctx, owner.DeviceID) if err == nil && currentPeer != peer { err = errors.New("Runtime connection changed during initialization") } @@ -153,16 +153,15 @@ func (w *Worker) prepareEnvironment(ctx context.Context, owner sessions.Environm if err == nil && index < len(cfg.Files) { var metadata environmentconfig.InitialFileMetadata var body []byte - metadata, body, err = w.dispatcher.SessionsReader.ReadInitialEnvironmentFile(step, owner.TenantID, owner.SessionID, index) + metadata, body, err = w.dispatcher.SessionsReader.ReadInitialEnvironmentFile(ctx, owner.TenantID, owner.SessionID, index) if err == nil { - err = installInitialFile(step, peer, identity, metadata, body) + err = installInitialFile(ctx, peer, identity, metadata, body) } } else if err == nil { command := operations[index-len(cfg.Files)] candidate = command.provisioningFailure(0) - err = runRuntimeSetup(step, peer, identity, command) + err = runRuntimeSetup(ctx, peer, identity, command) } - stop() if err != nil { var confirmed *runtimeStepFailure if errors.As(err, &confirmed) { diff --git a/services/core/internal/execution/runtime_setup.go b/services/core/internal/execution/runtime_setup.go index ccfcb795f..1f1822177 100644 --- a/services/core/internal/execution/runtime_setup.go +++ b/services/core/internal/execution/runtime_setup.go @@ -3,6 +3,7 @@ package execution import ( "context" "errors" + "time" "github.com/MiniMax-AI/OpenAgentCore/internal/agentcapabilities" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" @@ -85,7 +86,13 @@ func runRuntimeSetup(ctx context.Context, peer runtimePreparer, identity agentca if peer == nil { return errors.New("environment initialization request unavailable") } + deadline, bounded := ctx.Deadline() + budget := time.Until(deadline).Milliseconds() + if !bounded || budget <= 0 || budget > proto.RuntimePrepareMaxBudgetMS || ctx.Err() != nil { + return errors.New("environment initialization budget unavailable") + } request := operation.Request + request.BudgetMS = budget request.EnvironmentID, request.SessionID = identity.EnvironmentID, identity.SessionID result, err := peer.PrepareRuntime(ctx, uuid.NewString(), request, operation.Data) if err == nil && result.Outcome == "completed" { diff --git a/services/core/internal/execution/runtime_setup_test.go b/services/core/internal/execution/runtime_setup_test.go index be2a99461..2bd9be1bb 100644 --- a/services/core/internal/execution/runtime_setup_test.go +++ b/services/core/internal/execution/runtime_setup_test.go @@ -5,6 +5,7 @@ import ( "errors" "strings" "testing" + "time" "github.com/MiniMax-AI/OpenAgentCore/internal/agentcapabilities" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" @@ -36,7 +37,7 @@ func TestRuntimeSetupReceiptOutcomes(t *testing.T) { } { t.Run(test.name, func(t *testing.T) { peer := &receiptRuntime{result: proto.RuntimePrepareResultPayload{Outcome: test.outcome, ExitCode: test.code}, err: test.err} - err := runRuntimeSetup(t.Context(), peer, agentcapabilities.Identity{}, runtimeSetupOperation{Request: proto.RuntimePreparePayload{Action: "initialize", Initialization: &proto.RuntimeInitialization{Action: "setup", Command: setupCanary}}}) + err := runRuntimeSetup(initializationTestContext(t), peer, agentcapabilities.Identity{}, runtimeSetupOperation{Request: proto.RuntimePreparePayload{Action: "initialize", Initialization: &proto.RuntimeInitialization{Action: "setup", Command: setupCanary}}}) if test.outcome == "completed" && test.err == nil { if err != nil { t.Fatal(err) @@ -73,14 +74,47 @@ func TestInitialFileUsesTypedRuntimeBytes(t *testing.T) { size := int64(len(body)) owner := agentcapabilities.Identity{EnvironmentID: "environment", SessionID: "session"} peer := &receiptRuntime{result: proto.RuntimePrepareResultPayload{Outcome: "completed"}} - if err := installInitialFile(t.Context(), peer, owner, environmentconfig.InitialFileMetadata{Path: "/workspace/a", SizeBytes: &size}, body); err != nil { + if err := installInitialFile(initializationTestContext(t), peer, owner, environmentconfig.InitialFileMetadata{Path: "/workspace/a", SizeBytes: &size}, body); err != nil { t.Fatal(err) } if peer.request.Action != "file" || peer.request.File.Path != "/workspace/a" || peer.request.EnvironmentID != owner.EnvironmentID || peer.request.SessionID != owner.SessionID || string(peer.data) != setupCanary { t.Fatal("file transport changed") } size++ - if err := installInitialFile(t.Context(), peer, owner, environmentconfig.InitialFileMetadata{SizeBytes: &size}, body); err == nil { + if err := installInitialFile(initializationTestContext(t), peer, owner, environmentconfig.InitialFileMetadata{SizeBytes: &size}, body); err == nil { t.Fatal("mismatched source size accepted") } } + +func initializationTestContext(t *testing.T) context.Context { + t.Helper() + ctx, cancel := context.WithTimeout(t.Context(), 10*time.Minute) + t.Cleanup(cancel) + return ctx +} + +func TestRuntimeSetupUsesRemainingInitializationBudget(t *testing.T) { + ctx := initializationTestContext(t) + deadline, _ := ctx.Deadline() + peer := &receiptRuntime{result: proto.RuntimePrepareResultPayload{Outcome: "completed"}} + err := runRuntimeSetup(ctx, peer, agentcapabilities.Identity{}, runtimeSetupOperation{}) + if err != nil { + t.Fatal(err) + } + remaining := time.Until(deadline).Milliseconds() + if peer.request.BudgetMS < remaining || peer.request.BudgetMS > remaining+1000 || peer.request.BudgetMS <= 120000 { + t.Fatalf("remaining operation budget lost: %d vs %d", peer.request.BudgetMS, remaining) + } + for name, c := range map[string]context.Context{"unbounded": t.Context(), "expired": func() context.Context { + c, stop := context.WithDeadline(t.Context(), time.Now().Add(-time.Second)) + stop() + return c + }()} { + t.Run(name, func(t *testing.T) { + peer := &receiptRuntime{} + if runRuntimeSetup(c, peer, agentcapabilities.Identity{}, runtimeSetupOperation{}) == nil || peer.request.BudgetMS != 0 { + t.Fatal("invalid budget sent") + } + }) + } +} diff --git a/services/core/internal/runtimegateway/runtime_prepare.go b/services/core/internal/runtimegateway/runtime_prepare.go index 77e6493f9..677928f6a 100644 --- a/services/core/internal/runtimegateway/runtime_prepare.go +++ b/services/core/internal/runtimegateway/runtime_prepare.go @@ -50,8 +50,10 @@ func (s *Session) PrepareRuntime(ctx context.Context, id string, request proto.R s.capabilities = map[string]chan proto.Envelope{id: replies} s.capabilitiesMu.Unlock() defer func() { s.capabilitiesMu.Lock(); delete(s.capabilities, id); s.capabilitiesMu.Unlock() }() - ctx, cancel := context.WithTimeout(ctx, 195*time.Second) + ctx, cancel := context.WithTimeout(ctx, time.Duration(request.BudgetMS)*time.Millisecond) defer cancel() + transfer, stopTransfer := context.WithTimeout(ctx, time.Duration(proto.RuntimePrepareTransferBudgetMS)*time.Millisecond) + defer stopTransfer() exchange := func(payload proto.RuntimePreparePayload, outcome string, offset int) (proto.RuntimePrepareResultPayload, error) { env, err := proto.NewEnvelope(proto.TypeRuntimePrepare, id, payload) if err != nil { @@ -61,7 +63,13 @@ func (s *Session) PrepareRuntime(ctx context.Context, id string, request proto.R if err != nil || len(encoded) > proto.RuntimePrepareMaxFrameBytes { return unknown, errors.New("agentdaemon gateway: invalid Runtime frame") } - reply, err := s.exchangeChunkFrame(ctx, env, replies) + operation := ctx + if outcome != "completed" { + operation = transfer + } else if err := transfer.Err(); err != nil { + return unknown, err + } + reply, err := s.exchangeChunkFrame(operation, env, replies) if err != nil { return unknown, err } diff --git a/services/core/internal/runtimegateway/runtime_prepare_test.go b/services/core/internal/runtimegateway/runtime_prepare_test.go index 02834d87a..f0beea586 100644 --- a/services/core/internal/runtimegateway/runtime_prepare_test.go +++ b/services/core/internal/runtimegateway/runtime_prepare_test.go @@ -9,6 +9,7 @@ import ( "errors" "fmt" "testing" + "testing/synctest" "time" "github.com/MiniMax-AI/OpenAgentCore/internal/agentcapabilities" @@ -18,7 +19,7 @@ import ( ) func skillPreparation() proto.RuntimePreparePayload { - return proto.RuntimePreparePayload{EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Action: "skill", + return proto.RuntimePreparePayload{BudgetMS: 300000, EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Action: "skill", Skill: &agentskill.Metadata{Type: "inline", Name: "example", Description: "Example"}} } @@ -116,7 +117,7 @@ func TestCapabilitiesTransfersMoreThanFrameLimitAndCorrelates(t *testing.T) { func TestCapabilitiesFinalizeTransfersNoArchive(t *testing.T) { s := NewSession(newFakeConn(), "device", "tenant", "test", nil, nil) defer s.Close("test") - request := proto.RuntimePreparePayload{EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Action: "finalize", Sources: &agentcapabilities.Input{}} + request := proto.RuntimePreparePayload{BudgetMS: 300000, EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Action: "finalize", Sources: &agentcapabilities.Input{}} id := uuid.NewString() done := beginCapabilities(s, t.Context(), id, request, nil) env := nextCapabilityFrame(t, s) @@ -239,7 +240,7 @@ func TestRuntimeInitialFileChunking(t *testing.T) { t.Run(fmt.Sprint(size), func(t *testing.T) { s := NewSession(newFakeConn(), "device", "tenant", "test", nil, nil) defer s.Close("test") - request := proto.RuntimePreparePayload{EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Action: "file", File: &proto.RuntimeInitialFile{Path: "/workspace/project/file"}} + request := proto.RuntimePreparePayload{BudgetMS: 300000, EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Action: "file", File: &proto.RuntimeInitialFile{Path: "/workspace/project/file"}} data := bytes.Repeat([]byte("z"), size) id := uuid.NewString() done := beginCapabilities(s, t.Context(), id, request, data) @@ -289,7 +290,7 @@ func TestRuntimeInitializationNoDataAndExitReceipt(t *testing.T) { t.Run(fmt.Sprint(exit), func(t *testing.T) { s := NewSession(newFakeConn(), "device", "tenant", "test", nil, nil) defer s.Close("test") - request := proto.RuntimePreparePayload{EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Action: "initialize", Initialization: &proto.RuntimeInitialization{Action: "setup", Command: "echo test"}} + request := proto.RuntimePreparePayload{BudgetMS: 300000, EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Action: "initialize", Initialization: &proto.RuntimeInitialization{Action: "setup", Command: "echo test"}} if _, err := s.PrepareRuntime(t.Context(), uuid.NewString(), request, []byte("forbidden")); err == nil { t.Fatal("initialization body accepted") } @@ -327,3 +328,72 @@ func TestRuntimeInitializationNoDataAndExitReceipt(t *testing.T) { }) } } + +func TestRuntimePreparationSeparatesTransferAndApplyBudgets(t *testing.T) { + for _, commit := range []bool{false, true} { + t.Run(fmt.Sprint(commit), func(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + s := NewSession(newFakeConn(), "device", "tenant", "test", nil, nil) + defer s.Close("test") + request := proto.RuntimePreparePayload{BudgetMS: 300000, EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Action: "initialize", Initialization: &proto.RuntimeInitialization{Action: "configure"}} + id := uuid.NewString() + done := beginCapabilities(s, t.Context(), id, request, nil) + frame := nextCapabilityFrame(t, s) + var sent proto.RuntimePreparePayload + if frame.DecodePayload(&sent) != nil || sent.BudgetMS != request.BudgetMS { + t.Fatal("begin budget lost") + } + if commit { + replyCapabilities(s, id, proto.RuntimePrepareResultPayload{Outcome: "ready"}) + nextCapabilityFrame(t, s) + } + time.Sleep(121 * time.Second) + synctest.Wait() + if !commit { + result := finishCapabilities(t, done) + if !errors.Is(result.err, context.DeadlineExceeded) || result.result.Outcome != "unknown" { + t.Fatal(result) + } + return + } + select { + case result := <-done: + t.Fatal("transfer timer clamped execution", result) + default: + } + time.Sleep(180 * time.Second) + result := finishCapabilities(t, done) + if !errors.Is(result.err, context.DeadlineExceeded) || result.result.Outcome != "unknown" { + t.Fatal(result) + } + }) + }) + } +} + +func TestRuntimePreparationCommittedApplyStopsWaitingOnDisconnectOrCallerDeadline(t *testing.T) { + for _, disconnect := range []bool{false, true} { + t.Run(fmt.Sprint(disconnect), func(t *testing.T) { + s := NewSession(newFakeConn(), "device", "tenant", "test", nil, nil) + defer s.Close("test") + ctx, cancel := context.WithTimeout(t.Context(), 100*time.Millisecond) + defer cancel() + request := proto.RuntimePreparePayload{BudgetMS: 300000, EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Action: "initialize", Initialization: &proto.RuntimeInitialization{Action: "configure"}} + id := uuid.NewString() + done := beginCapabilities(s, ctx, id, request, nil) + nextCapabilityFrame(t, s) + replyCapabilities(s, id, proto.RuntimePrepareResultPayload{Outcome: "ready"}) + nextCapabilityFrame(t, s) + expected := error(context.DeadlineExceeded) + if disconnect { + s.Close("disconnect during apply") + expected = ErrSessionClosed + } + result := finishCapabilities(t, done) + if !errors.Is(result.err, expected) || result.result.Outcome != "unknown" { + t.Fatal(result) + } + noCapabilityFrame(t, s) + }) + } +}