diff --git a/deploy/kubernetes/BETA_CHANGELOG.md b/deploy/kubernetes/BETA_CHANGELOG.md index 424bd7a64..d5931817f 100644 --- a/deploy/kubernetes/BETA_CHANGELOG.md +++ b/deploy/kubernetes/BETA_CHANGELOG.md @@ -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` | +| 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. diff --git a/docs/sandbox-provider.md b/docs/sandbox-provider.md index 458bac041..41fc6e658 100644 --- a/docs/sandbox-provider.md +++ b/docs/sandbox-provider.md @@ -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. diff --git a/docs/zh/sandbox-provider.md b/docs/zh/sandbox-provider.md index f72a5f0b4..7c601db9a 100644 --- a/docs/zh/sandbox-provider.md +++ b/docs/zh/sandbox-provider.md @@ -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)。 @@ -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。 diff --git a/services/core/internal/execution/environment_admission.go b/services/core/internal/execution/environment_admission.go index 318fa2c34..5ed220472 100644 --- a/services/core/internal/execution/environment_admission.go +++ b/services/core/internal/execution/environment_admission.go @@ -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) diff --git a/services/core/internal/execution/runtime_fresh_hint_test.go b/services/core/internal/execution/runtime_fresh_hint_test.go new file mode 100644 index 000000000..6ca691379 --- /dev/null +++ b/services/core/internal/execution/runtime_fresh_hint_test.go @@ -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) + } + }) + } +} diff --git a/services/core/internal/execution/runtime_wake_hint.go b/services/core/internal/execution/runtime_wake_hint.go index ed7c9f1d6..b25fb4d03 100644 --- a/services/core/internal/execution/runtime_wake_hint.go +++ b/services/core/internal/execution/runtime_wake_hint.go @@ -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 { @@ -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 diff --git a/services/core/internal/execution/worker.go b/services/core/internal/execution/worker.go index 6783ddd1d..ae8053c67 100644 --- a/services/core/internal/execution/worker.go +++ b/services/core/internal/execution/worker.go @@ -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() @@ -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() diff --git a/services/core/internal/sandbox/e2b/create_observations.go b/services/core/internal/sandbox/e2b/create_observations.go new file mode 100644 index 000000000..cc65737a7 --- /dev/null +++ b/services/core/internal/sandbox/e2b/create_observations.go @@ -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) + } +} diff --git a/services/core/internal/sandbox/e2b/create_observations_test.go b/services/core/internal/sandbox/e2b/create_observations_test.go new file mode 100644 index 000000000..b5c996140 --- /dev/null +++ b/services/core/internal/sandbox/e2b/create_observations_test.go @@ -0,0 +1,94 @@ +package e2b + +import ( + "bytes" + "context" + "encoding/json" + "log/slog" + "os" + "strconv" + "strings" + "testing" + "time" + + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" + "github.com/google/uuid" +) + +func TestCreateObservationsRejectUnsafeDiagnostics(t *testing.T) { + valid := `{"event":"e2b_create_stage","stage":"sandbox_create","duration_us":1250,"completed":true}` + for _, bad := range []string{ + `private SDK secret`, + strings.Replace(valid, `"sandbox_create"`, `"private SDK secret"`, 1), + strings.Replace(valid, `1250`, `-1`, 1), + strings.Replace(valid, `1250`, `1800000001`, 1), + strings.Replace(valid, `1250`, `1.25`, 1), + strings.Replace(valid, `1250`, `null`, 1), + strings.Replace(valid, `true`, `null`, 1), + strings.Replace(valid, `true`, `"private SDK secret"`, 1), + strings.Replace(valid, `"completed":true`, `"credential":"private"`, 1), + strings.TrimSuffix(valid, "}") + `,"credential":"private"}`, + valid + `{}`, + strings.Repeat(" ", 256) + valid, + } { + if got := parseCreateObservations([]byte(bad)); len(got) != 0 { + t.Fatalf("unsafe diagnostic accepted: %q", bad) + } + } + got := parseCreateObservations([]byte("private SDK diagnostic\n" + valid + "\n" + valid)) + if len(got) != 1 || *got[0].DurationUS != 1250 || !*got[0].Completed { + t.Fatalf("valid timing lost or duplicated: %+v", got) + } +} + +func TestCreateObservationsRetainFailedStage(t *testing.T) { + got := parseCreateObservations([]byte(`{"event":"e2b_create_stage","stage":"bootstrap_run","duration_us":0,"completed":false}`)) + if len(got) != 1 || *got[0].Completed { + t.Fatal("failed stage discarded") + } +} + +func TestProcessCreateObservationsReachLogger(t *testing.T) { + for _, overflow := range []bool{false, true} { + t.Run(strconv.FormatBool(overflow), func(t *testing.T) { + p, caller, ref := fixture(t) + _, _ = p.Create(bounded(t), sandbox.Bootstrap{Reference: ref, SessionID: uuid.NewString(), DeviceID: uuid.NewString(), CoreURL: "https://core.example/api/v1", Credential: "private-bootstrap", Harness: "codex", NetworkAccess: "enabled"}) + q := caller.requests[0] + q.Deadline = time.Now().Add(time.Second) + raw := `{"event":"e2b_create_stage","stage":"sandbox_create","duration_us":1250,"completed":true}` + "\nprivate SDK secret\n" + if overflow { + raw += strings.Repeat("x", MaxOutputBytes) + } + // Keep arbitrary diagnostics in a fixture file, never shell source. + diagnostics := q.Config.StateDir + "/diagnostics" + if err := os.WriteFile(diagnostics, []byte(raw), 0600); err != nil { + t.Fatal(err) + } + script := "#!/bin/sh\ncat '" + diagnostics + "' >&2\nprintf '%s' '{\"Version\":" + strconv.Itoa(ProtocolVersion) + "}'\n" + if err := os.WriteFile(q.Config.Binary, []byte(script), 0700); err != nil { + t.Fatal(err) + } + var logs bytes.Buffer + previous := slog.Default() + slog.SetDefault(slog.New(slog.NewJSONHandler(&logs, nil))) + defer slog.SetDefault(previous) + _, err := (&ProcessCaller{}).Call(context.Background(), q) + if overflow { + if err == nil || logs.Len() != 0 { + t.Fatal("stderr overflow changed failure semantics or exported truncated diagnostics") + } + return + } + if err != nil { + t.Fatal(err) + } + var record map[string]any + if json.Unmarshal(logs.Bytes(), &record) != nil || record["allocation_id"] != ref.AllocationID || record["environment_id"] != ref.EnvironmentID || record["duration_ms"] != 1.25 || record["msg"] != "e2b create stage" { + t.Fatalf("stage did not reach correlated logger: %s", logs.Bytes()) + } + if strings.Contains(logs.String(), "private") { + t.Fatal("SDK diagnostic leaked") + } + }) + } +} diff --git a/services/core/internal/sandbox/e2b/process.go b/services/core/internal/sandbox/e2b/process.go index 1752c71ba..3e67cbded 100644 --- a/services/core/internal/sandbox/e2b/process.go +++ b/services/core/internal/sandbox/e2b/process.go @@ -54,6 +54,9 @@ func (p *ProcessCaller) Call(ctx context.Context, q Request) (Response, error) { case <-ctx.Done(): return Response{}, ctx.Err() case err = <-done: + if !stderr.exceeded { + observeHelperCreate(ctx, q, stderr.Bytes()) + } if err != nil || stdout.exceeded || stderr.exceeded { return Response{}, errors.New("helper result unconfirmed") } @@ -77,6 +80,12 @@ type limitBuffer struct { exceeded bool } +// Do not let the embedded bytes.Buffer's ReaderFrom bypass Write's limit when +// os/exec copies a child's output with io.Copy. +func (b *limitBuffer) ReadFrom(r io.Reader) (int64, error) { + return io.Copy(struct{ io.Writer }{b}, r) +} + func (b *limitBuffer) Write(data []byte) (int, error) { n := len(data) left := b.limit - b.Len() diff --git a/services/core/tools/e2b-provider/README.md b/services/core/tools/e2b-provider/README.md index 8e84b67e6..4233406de 100644 --- a/services/core/tools/e2b-provider/README.md +++ b/services/core/tools/e2b-provider/README.md @@ -56,6 +56,8 @@ Inspection uses SDK metadata and ID reads only. Inspection never uses SDK `conne ## Observations +Create emits six sequential timing records on helper stderr: `sandbox_create`, `ownership_check`, `template_check`, `bootstrap_write`, `bootstrap_run` and `ready_inspect`. Each record contains only `event=e2b_create_stage`, `stage`, integer `duration_us` and boolean `completed`. The Core adapter accepts only these stages, durations between zero and 1800 seconds, and one record per stage, then logs `e2b create stage` with the allocation and Environment IDs from its validated request. Unknown fields, oversized records and arbitrary SDK diagnostics are discarded. Timing uses a monotonic clock; `completed` describes the local step, not Runtime readiness. Records are exported when the helper exits before caller cancellation, including failed operations; the Core log timestamp is export time, not the time the step ran. Missing records are unobserved, not zero. Helper startup, lock waits and receipt work outside these steps remain part of the outer Provider duration. No request payload, credential, SDK error text or command output is logged. + `observe` reads each allocation's sandbox ID from its receipt without taking the allocation lock. It then runs one `GET /sandboxes/metrics` request and one labelled listing of the installation's running sandboxes concurrently, within the caller's deadline; the listing stops once every requested sandbox has appeared. Only a sandbox that the listing confirms for exactly that allocation is reported, with the listing's start time. A malformed metrics point makes only its row unavailable. Observation never connects to, renews or changes a sandbox and writes no receipt. [Runtime observability](../../../../contracts/agents-api/runtime-observability.md) owns the field mapping. ## Credential verification diff --git a/services/core/tools/e2b-provider/provider.py b/services/core/tools/e2b-provider/provider.py index 5a0fa2b70..241c4b934 100644 --- a/services/core/tools/e2b-provider/provider.py +++ b/services/core/tools/e2b-provider/provider.py @@ -1,8 +1,10 @@ """Five bounded SDK operations for an already authorized Core allocation, plus read-only deployment validation and batch observation.""" from concurrent.futures import ThreadPoolExecutor +from contextlib import contextmanager import json import math +import sys import time from datetime import datetime, timezone from uuid import UUID @@ -20,6 +22,24 @@ FIELDS = ('InstallationID', *REFERENCE_FIELDS) +@contextmanager +def create_stage(stage): + """Emit bounded numeric diagnostics, never SDK output or request contents.""" + started = time.monotonic_ns() + completed = False + try: + yield + completed = True + finally: + try: + sys.stderr.write(json.dumps({'event': 'e2b_create_stage', 'stage': stage, + 'duration_us': (time.monotonic_ns() - started) // 1000, + 'completed': completed}) + '\n') + except Exception: + # Observation must not change allocation or bootstrap outcomes. + pass + + def valid_id(value): try: return isinstance(value, str) and str(UUID(value)) == value and UUID(value).int != 0 @@ -213,14 +233,15 @@ def create(self): identity = dict(self.reference, InstallationID=self.config['InstallationID'], SessionID=bootstrap['SessionID'], DeviceID=bootstrap['DeviceID']) self.receipt.save(status='create_pending', bootstrap_identity=identity) - try: - cloud = Sandbox.create(template=self.config['Template'], timeout=self.config['TimeoutSeconds'], - metadata=self.metadata, lifecycle={'on_timeout': 'kill', 'auto_resume': False}, - **self.options()) - except Exception as error: - if definitely_rejected(error): - self.receipt.save(status='rejected', settled=True) - raise Failure('unconfirmed') from None + with create_stage('sandbox_create'): + try: + cloud = Sandbox.create(template=self.config['Template'], timeout=self.config['TimeoutSeconds'], + metadata=self.metadata, lifecycle={'on_timeout': 'kill', 'auto_resume': False}, + **self.options()) + except Exception as error: + if definitely_rejected(error): + self.receipt.save(status='rejected', settled=True) + raise Failure('unconfirmed') from None self.receipt.save(status='created', ids=[cloud.sandbox_id], connection=connection_material(cloud), compute={'Generation': 0, 'Name': self.reference['AllocationID'], 'ID': cloud.sandbox_id, 'RestoredFrom': None}) @@ -229,34 +250,39 @@ def create(self): # SDK Create returns connection material, but no metadata or resources. # Read its exact ID before writing credentials, even when Core adopts the # template's resources and does not supply explicit limits. - try: - detail = self.owns(Sandbox.get_info(cloud.sandbox_id, **self.options())) - self.qualified(detail) - except Failure: - self.receipt.save(status='configuration_rejected', settled=True) - raise + with create_stage('ownership_check'): + try: + detail = self.owns(Sandbox.get_info(cloud.sandbox_id, **self.options())) + self.qualified(detail) + except Failure: + self.receipt.save(status='configuration_rejected', settled=True) + raise # Validate the current template entry point before writing any credential. - check = run(cloud, {'Args': ['/usr/bin/python3', '-I', '-c', - "import os,sys; sys.exit(78 if not os.path.isfile('/opt/oac-e2b/managed_init.py') or not os.access('/opt/oac-e2b/managed_init.py', os.R_OK) else 0)"]}, - self.remaining, user='root') - if check['ExitCode'] != 0: - self.receipt.save(status='bootstrap_failed', settled=True) - raise Failure('template_invalid' if check['ExitCode'] == 78 else 'unconfirmed') + with create_stage('template_check'): + check = run(cloud, {'Args': ['/usr/bin/python3', '-I', '-c', + "import os,sys; sys.exit(78 if not os.path.isfile('/opt/oac-e2b/managed_init.py') or not os.access('/opt/oac-e2b/managed_init.py', os.R_OK) else 0)"]}, + self.remaining, user='root') + if check['ExitCode'] != 0: + self.receipt.save(status='bootstrap_failed', settled=True) + raise Failure('template_invalid' if check['ExitCode'] == 78 else 'unconfirmed') payload = dict(bootstrap, InstallationID=self.config['InstallationID'], RuntimeBootstrap=self.q['RuntimeBootstrap']) del payload['CoreURL'], payload['Credential'], payload['Harness'] if set(payload) != set(MANAGED_BOOTSTRAP_FIELDS): raise Failure('invalid') - cloud.files.write('/root/.oac/e2b/managed-bootstrap.json', json.dumps(payload), - user='root', request_timeout=self.remaining()) + with create_stage('bootstrap_write'): + cloud.files.write('/root/.oac/e2b/managed-bootstrap.json', json.dumps(payload), + user='root', request_timeout=self.remaining()) self.receipt.save(status='bootstrap_pending') - result = run(cloud, {'Args': ['/usr/bin/python3', '-I', '/opt/oac-e2b/managed_init.py']}, - self.remaining, user='root') - if result['ExitCode'] != 0: - self.receipt.save(status='bootstrap_failed', settled=True) - raise Failure('unconfirmed') + with create_stage('bootstrap_run'): + result = run(cloud, {'Args': ['/usr/bin/python3', '-I', '/opt/oac-e2b/managed_init.py']}, + self.remaining, user='root') + if result['ExitCode'] != 0: + self.receipt.save(status='bootstrap_failed', settled=True) + raise Failure('unconfirmed') self.receipt.save(status='bootstrap_exited', settled=True) - return self.inspect() + with create_stage('ready_inspect'): + return self.inspect() def renew(self): cloud = self.inspect() diff --git a/services/core/tools/e2b-provider/provider_test.py b/services/core/tools/e2b-provider/provider_test.py index 0c55dced7..feff8f1ba 100644 --- a/services/core/tools/e2b-provider/provider_test.py +++ b/services/core/tools/e2b-provider/provider_test.py @@ -1,5 +1,6 @@ """Controlled SDK boundary failures; live qualification remains separate.""" import copy +import io from datetime import datetime, timedelta, timezone import json from pathlib import Path @@ -12,7 +13,7 @@ from e2b import SandboxState from e2b.exceptions import AuthenticationException, SandboxNotFoundException -from provider import Provider +from provider import Provider, create_stage from helper_contract_generated import PROTOCOL_VERSION from sdk import restore, run from state import Failure, Receipt @@ -62,6 +63,46 @@ def call(self, operation): def record(self): return json.loads(next(Path(self.temporary.name).glob('*.json')).read_text()) + def test_create_stages_preserve_order_and_exclude_secrets(self): + stream = io.StringIO() + with patch('provider.sys.stderr', stream): + result = self.call('create') + self.assertEqual(result['ErrorCode'], '') + rows = [json.loads(line) for line in stream.getvalue().splitlines()] + self.assertEqual([row['stage'] for row in rows], [ + 'sandbox_create', 'ownership_check', 'template_check', + 'bootstrap_write', 'bootstrap_run', 'ready_inspect']) + for row in rows: + self.assertEqual(set(row), {'event', 'stage', 'duration_us', 'completed'}) + self.assertTrue(row['completed']) + self.assertGreaterEqual(row['duration_us'], 0) + self.assertNotIn('private', stream.getvalue()) + + def test_failed_create_stage_does_not_expose_exception_or_replay(self): + stream = io.StringIO() + self.api.create.side_effect = TimeoutError('private SDK secret') + with patch('provider.sys.stderr', stream): + self.assertEqual(self.call('create')['ErrorCode'], 'unconfirmed') + self.assertEqual(self.call('create')['ErrorCode'], 'exists') + rows = [json.loads(line) for line in stream.getvalue().splitlines()] + self.assertEqual(len(rows), 1) + self.assertEqual(rows[0]['stage'], 'sandbox_create') + self.assertFalse(rows[0]['completed']) + self.assertNotIn('private', stream.getvalue()) + self.api.create.assert_called_once() + + def test_broken_diagnostic_sink_does_not_change_create(self): + with patch('provider.sys.stderr') as stream: + stream.write.side_effect = OSError('sink closed') + self.assertEqual(self.call('create')['ErrorCode'], '') + + def test_stage_uses_monotonic_duration(self): + stream = io.StringIO() + with patch('provider.sys.stderr', stream), patch('provider.time.monotonic_ns', side_effect=[1000, 2501000]): + with create_stage('sandbox_create'): + pass + self.assertEqual(json.loads(stream.getvalue())['duration_us'], 2500) + def test_create_recover_and_never_replay(self): result = self.call('create') self.assertEqual(result['ErrorCode'], '')