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
13 changes: 13 additions & 0 deletions deploy/kubernetes/BETA_CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -95,3 +95,16 @@ The Environment file creation handler adds request- and trace-correlated static
Deploy through **core-deploy** from workflow branch `main`, selecting `source_branch=beta` and the full integrated commit already merged into beta, as described in the [Kubernetes guide](README.md#deploy). This replaces the existing production Core and Web; it does not create a separate beta environment. No Runtime template build or activation is required by this diagnostic change.

Verify both rollout revisions and ready Service endpoints, then check public health and authenticated API behavior. A safe rejected Environment file request can establish that the new event is correlated with its response IDs; it does not prove a successful upload or identify the cause of earlier 400 responses. This entry records integration requirements, not deployment success. Do not replay uncertain uploads or create paid execution solely to produce diagnostic evidence.

## 2026-10-08 — Fresh Environment provisioning hints and Create timings

| Item | Value |
| --- | --- |
| Fork beta baseline | `ccd9071b93008739d7ddc18c5b24eb45525a005c` |
| Integration branch | `codex/fresh-runtime-provision-hints` |

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Remove integration-branch process history

This row preserves a transient development branch rather than a current deployment requirement, so it becomes stale when that branch is renamed or deleted. Remove it and retain only current release and deployment facts; the repository documentation rules explicitly prohibit process history and task chronology.

AGENTS.md reference: AGENTS.md:L59-L59

Useful? React with 👍 / 👎.

| Schema migrations / SQL changes | None |
| Deployment requirements | Rebuild and deploy Core with its packaged E2B helper; no Runtime template change |

Committed hosted Environments and pending inputs notify their existing lifecycle owner. Hints preserve the lease, placement, capacity and one-shot allocation checks, with periodic recovery and one extra scan per maintenance period. The E2B helper adds bounded Create stage timings, and subprocess output collection enforces its byte limit through both copy paths. The [Sandbox Provider guide](../../docs/sandbox-provider.md#per-node-lifecycle-workers) and [E2B helper guide](../../services/core/tools/e2b-provider/README.md#observations) own the behavior and timing limitations.

Validation includes database-backed ordinary/stream creation and input recovery without advancing the maintenance clock, concurrent retry and lifecycle regressions, Go race checks, Python helper tests and independent review. Production timing must be verified after rollout; these checks do not establish an end-to-end latency guarantee.
2 changes: 2 additions & 0 deletions docs/sandbox-provider.md
Original file line number Diff line number Diff line change
Expand Up @@ -174,6 +174,8 @@ Terminal cleanup atomically revokes the device's authority, records the Environm

Each registered node has one serial lifecycle worker that owns its gate, allocation and pending cursors, connections and wake hints; E2B allocations share one serial lifecycle without a node. A thin coordinator discovers nodes and shuts workers down, and never holds its map mutex during database, provider or wait operations. Workers advance independently, so a stuck provider on one online node never stalls another: lifecycle concurrency is one operation per node and grows with the node count. Offline workers stay, so their retained resources remain observable after reconnection.

After a hosted Environment commits, Core sends a bounded hint to its placement's lifecycle worker. A committed pending input also sends a hint, recovering a missed creation notification. Hints use the existing serial gate, lease, capacity checks and one-shot allocation receipts; they do not provision inline in the admission handler. Each normal maintenance period permits at most one extra hinted scan, and periodic scans recover missed or coalesced hints. Existing allocation observations still precede pending provisioning, so a slow observation on the same node can delay a fresh Environment.

Allocation scans filter by node before their 32-row page limit, and pending scans join the unreleased committed placement. Each node advances its own cursor, including past failed observations, and wraps once at the end. Direct provisioning resolves the tenant-scoped placement before entering that node's gate, and an existing allocation must agree with it; Core never chooses another node.

Before releasing the execution lease, the coordinator stops accepting work and cancels and drains every node worker and direct caller. Lease loss affects everything; ordinary provider failures stay within their node. A planned deployment drain or node retirement cancels lifecycle contexts synchronously between leased operations, through the lease gate with its five-second bound, including an active manual reconcile, and never cancels an in-flight leased query just to change configuration. A failed cancellation fence closes manager admission and reports owner failure. A failed retirement keeps the original lifecycle identity and gate until owner shutdown, and the drain barrier stays closed. Session locks, deployment capacity transactions and revision-checked receipts stay authoritative, and no external operation holds a database lock.
Expand Down
4 changes: 3 additions & 1 deletion 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: e8bfcdfb5e0c445036753a64561903a28378e601dd0f7908cf9b387f22f2beaa
source_hash: 2597ae4c793c4607f8ffd3d264ad6c2dcabb13e473f9f7e595ac91319e16837c
---

**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 @@ -174,6 +174,8 @@ allocation、专用 daemon credential digest 和精确 Session binding 在 `Crea

### 每节点生命周期 worker {#per-node-lifecycle-workers}

托管 Environment 提交后,Core 会向其放置节点的生命周期 worker 发送有界提示。已提交的待处理输入也会发送提示,以恢复遗漏的创建通知。提示复用现有串行门控、执行租约、容量检查和一次性分配凭据,不会在准入处理器内直接创建沙箱。每个正常维护周期最多允许一次额外提示扫描;周期扫描负责恢复遗漏或合并的提示。已有分配的观察仍先于待创建资源处理,因此同一节点上的慢观察仍可能延迟新 Environment。

每个注册 node 有一个串行 lifecycle worker,负责 gate、allocation 与 pending cursor、connection 和 wake hint;E2B allocation 共享一个没有 node 的串行 lifecycle。薄 coordinator 发现 node 并关闭 worker,数据库、provider 或等待操作期间不持有 map mutex。worker 独立推进,因此一个在线 node 的 provider 卡住不会阻塞其他 node:lifecycle 并发为每 node 一项操作,随 node 数量增长。离线 worker 保留,因此其保留资源在重连后仍可观察。

allocation scan 在应用 32 行分页限制前按 node 过滤,pending scan join 尚未释放的已提交 placement。每个 node 推进自己的 cursor,包括越过失败观测,并在末尾回绕一次。direct provisioning 在进入该 node gate 前解析 tenant 范围内 placement,已有 allocation 必须与其一致;Core 不选择另一 node。
Expand Down
2 changes: 1 addition & 1 deletion services/core/internal/execution/environment_admission.go
Original file line number Diff line number Diff line change
Expand Up @@ -120,7 +120,7 @@ func (w *Worker) submitEnvironmentInputs(ctx context.Context, session sessions.S
"session_id", session.ID, "reservation_id", reservation.ID, "state", reservation.State)
}()
w.wakeScheduler()
if reservation.State == sessions.EnvironmentInputPending && !reservation.IsInitial {
if reservation.State == sessions.EnvironmentInputPending {
w.hintRuntimeWake(ctx, session)
}
ticker := time.NewTicker(250 * time.Millisecond)
Expand Down
193 changes: 193 additions & 0 deletions services/core/internal/execution/runtime_fresh_hint_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,193 @@
package execution

import (
"context"
"encoding/json"
"errors"
"sync"
"sync/atomic"
"testing"
"time"

v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/identity"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions"
"github.com/google/uuid"
)

type freshHintProvider struct {
sandbox.SandboxProvider
creates atomic.Int32
created chan struct{}
}

type hintSessionReader struct {
sessions.Reader
environment sessions.Environment
}

func (r hintSessionReader) GetSessionEnvironment(ctx context.Context, _, _ string) (sessions.Environment, error) {
return r.environment, ctx.Err()
}

func TestFreshHintRoutesAndPreservesLifecycleGuards(t *testing.T) {
for _, scenario := range []string{"direct", "node", "closed", "switching", "cancelled", "self_hosted", "released_placement", "lookup_failed", "suspended"} {
t.Run(scenario, func(t *testing.T) {
m := testRuntimeManager(t)
m.config = RuntimeProvider{InstallationID: uuid.NewString()}
nodeID := ""
if scenario == "node" || scenario == "released_placement" {
m.config.ProviderKind = "docker"
nodeID = uuid.NewString()
}
node, err := m.node(nodeID)
if err != nil {
t.Fatal(err)
}
environment := sessions.Environment{ID: uuid.NewString(), Configuration: json.RawMessage(`{"type":"openai_hosted"}`)}
if scenario == "self_hosted" {
environment.Configuration = json.RawMessage(`{"type":"self_hosted","workspace_directory":"/workspace"}`)
}
if scenario == "suspended" {
environment.Initialization = "complete"
m.config.Suspension = &RuntimeSuspensionPolicy{}
}
reader := &strictDeploymentReader{t: t,
environmentAllocation: func(context.Context, deployment.AllocationKey) (deployment.Allocation, error) {
if scenario == "lookup_failed" {
return deployment.Allocation{}, errors.New("database unavailable")
}
if scenario == "suspended" {
return deployment.Allocation{ProviderKey: m.config.InstallationID, State: "running", CreateSettled: true, ComputePhase: "suspended"}, nil
}
return deployment.Allocation{}, deployment.ErrNotFound
},
lifecyclePlacement: func(context.Context, deployment.AllocationKey) (deployment.LifecyclePlacement, error) {
return deployment.LifecyclePlacement{Provider: m.config.ProviderKind, PlacementNodeID: nodeID, PlacementReleased: scenario == "released_placement"}, nil
},
}
m.deploymentService, _ = deploymentOperations(t, &strictDeploymentStorage{t: t}, reader, &strictExecutionStorage{t: t})
worker := &Worker{runtimes: m, dispatcher: &Dispatcher{SessionsReader: hintSessionReader{environment: environment}, DeploymentReader: reader}}
ctx, cancel := context.WithCancel(t.Context())
defer cancel()
switch scenario {
case "closed":
m.closed = true
case "switching":
m.switching = true
case "cancelled":
cancel()
}
worker.hintRuntimeWake(ctx, sessions.Session{TenantID: uuid.NewString(), ID: uuid.NewString()})
want := 0
if scenario == "direct" || scenario == "node" || scenario == "suspended" {
want = 1
}
if got := len(node.lifecycle.wakeHints); got != want {
t.Fatalf("hints = %d, want %d", got, want)
}
})
}
}

func (p *freshHintProvider) Create(_ context.Context, b sandbox.Bootstrap) (sandbox.Info, error) {
p.creates.Add(1)
select {
case p.created <- struct{}{}:
default:
}
return sandbox.Info{Reference: b.Reference, ProviderID: b.AllocationID, State: "running", BootstrapComplete: true, CreateSettled: true}, nil
}

func TestFreshEnvironmentHintProvisionsWithoutMaintenanceTick(t *testing.T) {
for _, mode := range []string{"create", "stream", "recovered_input", "recovered_initial"} {
t.Run(mode, func(t *testing.T) {
s, owner, deployments, reader, pool := resetManagerStoreDB(t, nil)
installation := initializeE2BDeployment(t, owner)
sessionReader, _ := testSessions(t, pool, testCredentialCipher(t))
provider := &freshHintProvider{created: make(chan struct{}, 1)}
m := testRuntimeManager(t)
m.store, m.sessions, m.sessionExecution = owner.Store, sessionReader, owner.Sessions
m.deployment, m.deploymentService, m.deploymentReader = owner.Deployment, deployments, reader
m.lease, m.registry = owner.Lease, runtimegateway.NewRegistry()
m.config = RuntimeProvider{InstallationID: installation, CoreURL: fixturePublicURL + "/api/v1", Provider: provider}
worker := &Worker{admission: s, lease: owner.Lease, runtimes: m, dispatcher: &Dispatcher{SessionsReader: sessionReader, DeploymentReader: reader, notifications: &executionNotifications{}}, scheduleWake: make(chan struct{}, 1)}
node, err := m.node("")
if err != nil {
t.Fatal(err)
}
ctx, cancel := context.WithCancel(t.Context())
done, scans := make(chan error, 1), make(chan struct{}, 4)
// No value is ever sent to ticks: provisioning must come from a hint.
ticks := make(chan time.Time)
go func() {
done <- runRuntimeMaintenance(ctx, ticks, node.lifecycle.wakeHints, func(ctx context.Context) error {
err := node.lifecycle.reconcile(ctx)
scans <- struct{}{}
return err
})
}()
t.Cleanup(func() { cancel(); <-done })
<-scans // Startup scan finishes before any Session is committed.
tenant := uuid.NewString()
input := sessions.CreateSession{Creator: identity.Subject{Kind: "service_account", ID: "fixture"}, Engine: "codex", IdempotencyKey: uuid.NewString(), Configuration: json.RawMessage(`{"agent":{"model":"test-model"},"environment":{"type":"openai_hosted"}}`), ModelProvider: &v1.ModelProviderInput{Protocol: "responses", BaseURL: "https://model.fixture.example/v1", APIKey: "fixture-key"}, ModelProviderSource: v1.ModelProviderSourceSession}
messages := []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"input":[{"role":"user","content":[{"type":"input_text","text":"fixture"}]}]}`)}}
if mode == "recovered_initial" {
input.InitialInputs = messages
}
var session sessions.Session
switch mode {
case "create":
session, err = worker.CreateSession(ctx, tenant, input)
case "stream":
var creation sessions.Creation
creation, err = worker.CreateSessionStream(ctx, tenant, input)
session = creation.Session
case "recovered_input", "recovered_initial":
session, err = s.CreateSession(ctx, tenant, input)
}
if err != nil {
t.Fatal(err)
}
if mode == "recovered_input" || mode == "recovered_initial" {
key := uuid.NewString()
if mode == "recovered_initial" {
if err := pool.QueryRow(ctx, "SELECT idempotency_key FROM environment_input_reservations WHERE session_id=$1 AND is_initial", session.ID).Scan(&key); err != nil {
t.Fatal(err)
}
}
inputDone := make(chan error, 20)
inputCtx, stopInput := context.WithCancel(ctx)
for range 20 {
go func() {
_, inputErr := worker.submitEnvironmentInputs(inputCtx, session, key, messages)
inputDone <- inputErr
}()
}
t.Cleanup(func() {
stopInput()
for range 20 {
<-inputDone
}
})
}
select {
case <-provider.created:
case <-time.After(2 * time.Second):
t.Fatal("committed Environment waited for a maintenance tick")
}
var requests sync.WaitGroup
for range 20 {
requests.Go(func() { worker.hintRuntimeWake(ctx, session) })
}
requests.Wait()
<-scans // Wait for the hinted scan to settle its durable allocation.
if got := provider.creates.Load(); got != 1 {
t.Fatalf("concurrent hints issued %d Creates", got)
}
})
}
}
33 changes: 29 additions & 4 deletions services/core/internal/execution/runtime_wake_hint.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,14 +2,15 @@ package execution

import (
"context"
"errors"
"time"

"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions"
)

// A hint only accelerates observation of an already committed input. Lookup or
// delivery failure leaves that input for the normal maintenance scan.
// Hints only accelerate lifecycle work for committed Environments and inputs.
// Lookup or delivery failure leaves that work for the normal maintenance scan.
func (w *Worker) hintRuntimeWake(ctx context.Context, session sessions.Session) {
r := w.runtimes
if r == nil {
Expand All @@ -19,16 +20,40 @@ func (w *Worker) hintRuntimeWake(ctx context.Context, session sessions.Session)
config := r.config
available := !r.closed && !r.switching
r.mu.Unlock()
if !available || config.Suspension == nil || r.ctx.Err() != nil {
if !available || r.ctx.Err() != nil {
return
}
lookup, cancel := context.WithTimeout(ctx, time.Second)
defer cancel()
environment, err := w.dispatcher.SessionsReader.GetSessionEnvironment(lookup, session.TenantID, session.ID)
if err != nil || environment.Initialization != "complete" {
if err != nil {
return
}
placement, err := parseEnvironmentPlacement(environment.Configuration)
if err != nil || placement.Type != "openai_hosted" {
return
}
owner, err := w.dispatcher.DeploymentReader.EnvironmentAllocation(lookup, deployment.AllocationKey{TenantID: session.TenantID, EnvironmentID: environment.ID})
if errors.Is(err, deployment.ErrNotFound) {
// Placement exists before allocation. Route through its lifecycle owner;
// the scan still enforces lease, capacity and one-shot Create receipts.
nodeID, err := r.deploymentService.LifecycleNode(lookup, session.TenantID, environment.ID)
if err != nil || lookup.Err() != nil {
return
}
node, err := r.node(nodeID)
if err != nil {
return
}
select {
case node.lifecycle.wakeHints <- struct{}{}:
default:
}
return
}
if config.Suspension == nil || environment.Initialization != "complete" {
return
}
if err != nil || owner.ProviderKey != config.InstallationID || owner.State != "running" ||
!owner.CreateSettled || owner.SessionDeleted || owner.Expired {
return
Expand Down
6 changes: 6 additions & 0 deletions services/core/internal/execution/worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -166,6 +166,9 @@ func (w *Worker) CreateSession(ctx context.Context, tenant string, input session
return sessions.Session{}, err
}
session, err := w.admission.CreateSession(ctx, tenant, input)
if err == nil {
w.hintRuntimeWake(ctx, session)
}
if err == nil && len(input.InitialInputs) > 0 {
recordInitialInputOrigin(ctx, session.ID)
w.wakeScheduler()
Expand All @@ -179,6 +182,9 @@ func (w *Worker) CreateSessionStream(ctx context.Context, tenant string, input s
return sessions.Creation{}, err
}
creation, err := w.admission.CreateSessionStream(ctx, tenant, input)
if err == nil {
w.hintRuntimeWake(ctx, creation.Session)
}
if err == nil && len(input.InitialInputs) > 0 {
recordInitialInputOrigin(ctx, creation.Session.ID)
w.wakeScheduler()
Expand Down
56 changes: 56 additions & 0 deletions services/core/internal/sandbox/e2b/create_observations.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
package e2b

import (
"bytes"
"context"
"encoding/json"
"io"

obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log"
)

type createObservation struct {
Event string `json:"event"`
Stage string `json:"stage"`
DurationUS *int64 `json:"duration_us"`
Completed *bool `json:"completed"`
}

// Helper stderr is untrusted. Only fixed stage names and bounded numeric fields
// reach the logger; arbitrary SDK diagnostics are never forwarded.
func parseCreateObservations(raw []byte) []createObservation {
var out []createObservation
seen := make(map[string]bool)
for _, line := range bytes.Split(raw, []byte{'\n'}) {
if len(line) > 256 {
continue
}
var value createObservation
decoder := json.NewDecoder(bytes.NewReader(line))
decoder.DisallowUnknownFields()
if decoder.Decode(&value) != nil || value.Event != "e2b_create_stage" || value.DurationUS == nil || value.Completed == nil {
continue
}
var extra any
if decoder.Decode(&extra) != io.EOF || *value.DurationUS < 0 || *value.DurationUS > 1_800_000_000 || seen[value.Stage] {
continue
}
switch value.Stage {
case "sandbox_create", "ownership_check", "template_check", "bootstrap_write", "bootstrap_run", "ready_inspect":
seen[value.Stage] = true
out = append(out, value)
}
}
return out
}

func observeHelperCreate(ctx context.Context, q Request, raw []byte) {
if q.Operation != "create" {
return
}
for _, value := range parseCreateObservations(raw) {
obslog.Info(ctx, "e2b create stage", "stage", value.Stage,
"duration_ms", float64(*value.DurationUS)/1000, "completed", *value.Completed,
"environment_id", q.Reference.EnvironmentID, "allocation_id", q.Reference.AllocationID)
}
}
Loading
Loading