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
2 changes: 1 addition & 1 deletion docs/sandbox-provider.md
Original file line number Diff line number Diff line change
Expand Up @@ -196,7 +196,7 @@ Queued work and live Environment file access wake a suspended Environment; histo

Snapshot retention bounds compute resources, not the lifetime of an eligible retained Session. After the deadline, Core can clean up the allocation without terminating the Session only when its independent filesystem object is ready, initialization is complete, the Session and Environment remain nonterminal, and the previous authenticated Runtime has a persisted `retained_native_history` declaration. The capability and native recovery obligations are defined by the [Core–Runtime protocol](runtime-protocol.md). Filesystem retention alone is insufficient.

When compatible demand cannot obtain a retained slot, Core may end an idle, stopped allocation’s checkpoint retention early. It rechecks native-history eligibility, pending work and capacity under the Session and deployment locks. The same expiry cleanup must confirm resource deletion before the waiting demand can reserve that slot. The 24-hour retention is an upper bound, not a guarantee of memory preservation under capacity pressure; files and qualified native history remain retained.
Capacity pressure never shortens checkpoint retention. Until its configured deadline, a suspended allocation keeps its snapshot and retained slot; unavailable capacity or an incompatible node does not authorize cold replacement. Pending requests retain their original input deadline. After retention expires, the ordinary expiry cleanup must confirm resource deletion before another allocation can reserve that slot; files and qualified native history remain retained.

Core confirms the original creation is settled and all owned compute and snapshots are cleaned up before releasing allocation and placement ownership. The old device's authority is revoked before replacement admission. Until Create, Kill or DeleteSnapshot has a confirmed outcome, it keeps ownership and admits no replacement writer. A recoverable Session and Environment remain disconnected with the same filesystem object and initialization result. Archive, reset and deletion take precedence and cannot be undone by a wake request.

Expand Down
4 changes: 2 additions & 2 deletions docs/zh/sandbox-provider.md
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
---
title: "添加 Sandbox Provider"
source: docs/sandbox-provider.md
source_hash: 4a7ce15e8fc2c9519b4974a60c746fbd436b24bbce5ab2ef461e442bea3e5f86
source_hash: 567cd130134a45f483b156d7f2b74d246b191dd5db77b976976c0d710414dd37
---

**Sandbox Provider** 为 Core 管理的 Environment 提供 Runtime daemon 运行所需的外层计算资源,以及启动 daemon 的有界引导流程。本指南说明如何添加 Provider,并作为 Core 驱动 Provider 的参考。接口为 [`SandboxProvider`](https://github.com/MiniMax-AI/OpenAgentCore/blob/main/services/core/internal/sandbox/sandbox_provider.go)。
Expand Down Expand Up @@ -198,7 +198,7 @@ Worker lease、Session lock 与 per-node gate 对每个 provider 负责 suspensi

快照保留期限制计算资源,不限制符合条件的保留 Session 的生命周期。到期后,只有独立文件系统对象已就绪、初始化已完成、Session 与 Environment 均未终结,且此前已认证 Runtime 的 `retained_native_history` 声明已持久化时,Core 才能在不终结 Session 的情况下清理 allocation。能力及原生恢复义务由 [Core–Runtime 协议](runtime-protocol.md)定义。仅保留文件系统并不足够。

兼容需求无法获得 retained slot 时,Core 可以提前结束空闲且已停止 allocation 的检查点保留期。它在 Session 和 deployment 锁内重新检查原生历史资格、pending work 和容量。等待需求预留该 slot 前,同一 expiry cleanup 必须确认资源已删除。24 小时保留期是上限,不保证容量压力下仍保留内存;文件与已验证的原生历史继续保留。
容量压力不会缩短检查点保留期。在配置的期限到达前,已挂起的 allocation 保留快照并继续占用 retained slot;容量不足或节点不兼容不允许降级为冷替换。等待请求保持原有的输入期限。保留期到期后,常规 expiry cleanup 必须确认资源已删除,其他 allocation 才能预留该 slot;文件与已验证的原生历史继续保留。

Core 在释放 allocation 和 placement 归属前,确认原始创建已结算且所属计算资源与快照均已清理。替代计算资源准入前撤销旧 device 的权限。Create、Kill 或 DeleteSnapshot 的结果尚未确认时,Core 保留归属,不准入替代写入方。可恢复的 Session 与 Environment 保持断连,保留同一文件系统对象和初始化结果。归档、重置和删除优先,唤醒请求不能撤销它们。

Expand Down
44 changes: 0 additions & 44 deletions services/core/internal/deployment/retention_pressure.go

This file was deleted.

23 changes: 7 additions & 16 deletions services/core/internal/deployment/suspension_pressure.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,20 +40,16 @@ func (e *ExecutionOperations) canSuspendForDemand(tx AllocationTx, owner Allocat
if pressure.RestoreWaiting {
return !slices.Contains(pressure.InFlightNodes, owner.NodeID), nil
}
reclaimRetained := source.Retained >= int64(source.MaxRetained)
if reclaimRetained {
retained, err := tx.CanRetainEnvironment(owner)
if err != nil || !retained {
return false, err
}
if source.Retained >= int64(source.MaxRetained) {
return false, nil
}
return e.placementNeedsCapacity(tx, pressure, owner.NodeID, true, reclaimRetained)
return e.placementNeedsCapacity(tx, pressure, owner.NodeID)
}

// placementNeedsCapacity scans bounded pages under the caller's deployment
// lock and transaction deadline. Suspension and early retention expiry use the
// same admission decision; an ineligible old receipt never supplies pressure.
func (e *ExecutionOperations) placementNeedsCapacity(tx AllocationTx, pressure SuspensionDemand, nodeID string, freeActive, freeRetained bool) (bool, error) {
// lock and transaction deadline. Suspension must make active capacity usable
// without releasing a retained slot; an ineligible receipt supplies no pressure.
func (e *ExecutionOperations) placementNeedsCapacity(tx AllocationTx, pressure SuspensionDemand, nodeID string) (bool, error) {
var after PlacementDemandCursor
for {
demands, next, err := tx.PlacementDemand(after)
Expand All @@ -80,12 +76,7 @@ func (e *ExecutionOperations) placementNeedsCapacity(tx AllocationTx, pressure S
}
for i := range nodes {
if nodes[i].ID == nodeID {
if freeActive {
nodes[i].Active--
}
if freeRetained && nodes[i].Retained >= int64(nodes[i].MaxRetained) {
nodes[i].Retained--
}
nodes[i].Active--
}
}
chosen, err := e.service.rules.DecidePlacement(pressure.Deployment, nodes)
Expand Down
3 changes: 1 addition & 2 deletions services/core/internal/execution/runtime_compute.go
Original file line number Diff line number Diff line change
Expand Up @@ -221,8 +221,7 @@ func (r *runtimeLifecycle) restoreIdleCompute(ctx context.Context, p sandbox.San
return err
}
if !activity.Busy && !activity.WakeRequested {
_, err := r.deployment.EndRetentionForDemand(ctx, owner)
return err
return nil
}
if state.Snapshot == nil || state.Target != nil {
return sandbox.ErrOwnership
Expand Down
58 changes: 58 additions & 0 deletions services/core/internal/execution/runtime_replacement_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,10 +11,13 @@ import (
"testing"
"time"

v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1"
"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto"
"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto/prototest"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/adminaudit"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment/placement"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/identity"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/providercontract"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox"
Expand Down Expand Up @@ -51,6 +54,61 @@ func (p *retentionProvider) DeleteSnapshot(context.Context, sandbox.Reference, s
return nil
}

func TestCapacityPressurePreservesSuspendedSnapshot(t *testing.T) {
provider := &retentionProvider{}
fixture := newWorkspaceSettlementFixture(t, provider)
r, session := fixture.lifecycle, fixture.session
owner, err := r.provision(t.Context(), session.TenantID, session.Environment.ID, r.config.InstallationID)
if err != nil {
t.Fatal(err)
}
for _, statement := range []struct {
sql string
arg any
}{
{`UPDATE runtime_nodes SET max_active=1,max_retained=1 WHERE id=$1`, owner.NodeID},
{`UPDATE environments SET initialization='complete',status='disconnected' WHERE id=$1`, owner.EnvironmentID},
{`UPDATE devices SET supported_agent_kinds='[{"kind":"codex","available":true,"capabilities":{"retained_native_history":true}}]' WHERE id=$1`, owner.DeviceID},
{`UPDATE runtime_allocations SET state='running',compute_phase='suspended',compute_state='{"current":{"id":"old-compute"},"snapshot":{"id":"old-snapshot"}}',compute_retained_until=clock_timestamp()+interval '24 hours' WHERE id=$1`, owner.ID},
} {
if _, err := fixture.pool.Exec(t.Context(), statement.sql, statement.arg); err != nil {
t.Fatal(err)
}
}
waiting, err := fixture.sessions.CreateSession(t.Context(), session.TenantID, sessions.CreateSession{SupportsRetainedNativeHistory: true, Creator: identity.Subject{Kind: "service_account", ID: "fixture"}, Engine: "codex", IdempotencyKey: uuid.NewString(), Configuration: json.RawMessage(`{"agent":{"model":"test-model"},"environment":{"type":"openai_hosted","network":{"access":"disabled"}}}`), ModelProvider: &v1.ModelProviderInput{Protocol: "responses", BaseURL: "https://model.fixture.example/v1", APIKey: "fixture-key"}, ModelProviderSource: v1.ExecutionSourceSession})
if err != nil {
t.Fatal(err)
}
input, err := fixture.sessions.ReserveEnvironmentInput(t.Context(), session.TenantID, waiting.Session.ID, "waiting-for-capacity", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"input":[{"role":"user","content":[{"type":"input_text","text":"continue"}]}]}`)}})
if err != nil {
t.Fatal(err)
}
owner, err = r.reader.EnvironmentAllocation(t.Context(), owner.Key())
if err != nil {
t.Fatal(err)
}
for range 3 {
if err := r.observeCompute(t.Context(), owner); err != nil {
t.Fatal(err)
}
current, err := r.reader.EnvironmentAllocation(t.Context(), owner.Key())
if err != nil || current.Expired || current.ComputePhase != "suspended" || current.ComputeRevision != owner.ComputeRevision || current.ComputeRetainedUntil == nil || !current.ComputeRetainedUntil.Equal(*owner.ComputeRetainedUntil) {
t.Fatal("pressure changed retained snapshot", current, err)
}
if _, err := r.deployment.EnsurePlacement(t.Context(), deployment.AllocationKey{TenantID: session.TenantID, EnvironmentID: waiting.Session.Environment.ID}, r.config.InstallationID); !errors.Is(err, placement.ErrNodeUnavailable) {
t.Fatal("released retained capacity without cleanup", err)
}
}
pending, err := r.sessions.GetEnvironmentInputReservation(t.Context(), session.TenantID, waiting.Session.ID, input.ID)
if err != nil || pending.State != sessions.EnvironmentInputPending || !pending.Deadline.Equal(input.Deadline) {
t.Fatal("pressure changed input deadline", pending, err)
}

if provider.kills != 0 || provider.snapshots != 0 || provider.creates != 1 {
t.Fatal("pressure changed native resources", provider)
}
}

func TestExpiredRetainedComputePreservesSessionAndPendingInput(t *testing.T) {
provider := &retentionProvider{killError: sandbox.ErrComputeUnconfirmed}
fixture := newWorkspaceSettlementFixture(t, provider)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,12 @@ package deploymentpg_test

import (
"encoding/json"
"errors"
"testing"
"time"

"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment/placement"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox"
"github.com/google/uuid"
)
Expand Down Expand Up @@ -37,26 +39,36 @@ func TestRetentionPressureKeepsHistoryUntilOrdinaryCleanupSettles(t *testing.T)
f := newFixture(t)
changes, _ := f.execution(t)
installation, view := f.initialize(t, changes, sandbox.Selection{Provider: "microsandbox", DeploymentSpec: retainedSpecification()})
node := f.enroll(t, view, deployment.Capacity{MaxActive: 1, MaxRetained: 1})
node := f.enroll(t, view, deployment.Capacity{MaxActive: 3, MaxRetained: 3})
f.connect(t, node.NodeID)
owner := qualifyPressureRetention(t, f, runningPressureOwner(t, f, changes, installation, node.NodeID, view.Generation))
for range 2 {
qualifyPressureRetention(t, f, runningPressureOwner(t, f, changes, installation, node.NodeID, view.Generation))
}
waiting := pendingPressureDemand(t, f)
ended, err := changes.EndRetentionForDemand(t.Context(), owner)
if err != nil {
t.Fatal(err)
for range 6 {
pendingPressureDemand(t, f)
}
if ended.ComputeRetainedUntil == nil || !ended.ComputeRetainedUntil.Before(*owner.ComputeRetainedUntil) {
t.Fatal("retention was not ended", ended)
if _, err := changes.EnsurePlacement(t.Context(), waiting, installation); !errors.Is(err, placement.ErrNodeUnavailable) {
t.Fatal("reserved despite retained capacity limit")
}
current, err := f.adapter.EnvironmentAllocation(t.Context(), owner.Key())
if err != nil || current.Expired || !current.ComputeRetainedUntil.Equal(*owner.ComputeRetainedUntil) {
t.Fatal("pressure shortened retention", current, err)
}
// Advance the fixture clock past the authored deadline, not because of demand.
if _, err := f.pool.Exec(t.Context(), `UPDATE runtime_allocations SET compute_retained_until=clock_timestamp()-interval '1 second' WHERE id=$1`, owner.ID); err != nil {
t.Fatal(err)
}
// The database deadline is only an intent. It is not a cleanup receipt.
nodes, err := f.adapter.Nodes(t.Context())
if err != nil || nodes[0].Retained != 1 {
if err != nil || nodes[0].Retained != 3 || nodes[0].Active != 0 {
t.Fatal("released capacity before cleanup", nodes, err)
}
if _, err = changes.EnsurePlacement(t.Context(), waiting, installation); err == nil {
if _, err = changes.EnsurePlacement(t.Context(), waiting, installation); !errors.Is(err, placement.ErrNodeUnavailable) {
t.Fatal("reserved before cleanup")
}
current, err := f.adapter.EnvironmentAllocation(t.Context(), owner.Key())
current, err = f.adapter.EnvironmentAllocation(t.Context(), owner.Key())
if err != nil || !current.Expired {
t.Fatal("database expiry", current, err)
}
Expand All @@ -78,50 +90,7 @@ func TestRetentionPressureKeepsHistoryUntilOrdinaryCleanupSettles(t *testing.T)
}
}

func TestRetentionPressureGuards(t *testing.T) {
for _, name := range []string{"no demand", "retained headroom", "wake requested", "missing native history", "free other node", "already ending"} {
t.Run(name, func(t *testing.T) {
f := newFixture(t)
changes, _ := f.execution(t)
installation, view := f.initialize(t, changes, sandbox.Selection{Provider: "microsandbox", DeploymentSpec: retainedSpecification()})
limit := 1
if name == "retained headroom" || name == "already ending" {
limit = 2
}
node := f.enroll(t, view, deployment.Capacity{MaxActive: 1, MaxRetained: limit})
f.connect(t, node.NodeID)
owner := qualifyPressureRetention(t, f, runningPressureOwner(t, f, changes, installation, node.NodeID, view.Generation))
if name != "no demand" {
pendingPressureDemand(t, f)
}
var err error
switch name {
case "wake requested":
_, err = f.pool.Exec(t.Context(), `UPDATE runtime_allocations SET compute_wake_requested=true WHERE id=$1`, owner.ID)
case "missing native history":
_, err = f.pool.Exec(t.Context(), `UPDATE devices SET supported_agent_kinds='[]' WHERE id=$1`, owner.DeviceID)
case "free other node":
other := f.enroll(t, view, deployment.Capacity{MaxActive: 1, MaxRetained: 2})
f.connect(t, other.NodeID)
case "already ending":
other := qualifyPressureRetention(t, f, runningPressureOwner(t, f, changes, installation, node.NodeID, view.Generation))
_, err = f.pool.Exec(t.Context(), `UPDATE runtime_allocations SET compute_retained_until=clock_timestamp()-interval '1 second' WHERE id=$1`, other.ID)
}
if err != nil {
t.Fatal(err)
}
next, err := changes.EndRetentionForDemand(t.Context(), owner)
if err != nil {
t.Fatal(err)
}
if next.ComputeRevision != owner.ComputeRevision || next.ComputeRetainedUntil == nil || !next.ComputeRetainedUntil.Equal(*owner.ComputeRetainedUntil) {
t.Fatal("ended retention despite guard", next)
}
})
}
}

func TestPressureSuspendsRetainableOwnerWhenRetainedSlotsAreFull(t *testing.T) {
func TestPressurePreservesRetainableOwnerWhenRetainedSlotsAreFull(t *testing.T) {
f := newFixture(t)
changes, _ := f.execution(t)
installation, view := f.initialize(t, changes, sandbox.Selection{Provider: "microsandbox", DeploymentSpec: retainedSpecification()})
Expand All @@ -137,7 +106,7 @@ func TestPressureSuspendsRetainableOwnerWhenRetainedSlotsAreFull(t *testing.T) {
}
pendingPressureDemand(t, f)
until := time.Now().Add(24 * time.Hour)
if _, err = changes.SetCompute(t.Context(), owner, "quiescing", json.RawMessage(`{}`), &until, 5*time.Minute); err != nil {
t.Fatal(err)
if _, err = changes.SetCompute(t.Context(), owner, "quiescing", json.RawMessage(`{}`), &until, 5*time.Minute); !errors.Is(err, deployment.ErrNotIdle) {
t.Fatal("suspended without usable capacity", err)
}
}
Loading
Loading