From 234ed63e9c9120d84d9ca9db55c596dd3fcbfb00 Mon Sep 17 00:00:00 2001 From: sam Date: Wed, 30 Sep 2026 09:32:08 +0800 Subject: [PATCH 01/14] Pause idle E2B sandboxes and resume original compute --- apps/docs/content/docs/sandbox-provider.mdx | 12 +- apps/docs/content/guide-sources.json | 4 +- contracts/agents-api/sandbox-deployment.md | 21 +- docs/sandbox-provider.md | 12 +- .../agents-api/cmd/server/managed_setup.go | 4 +- .../internal/execution/runtime_compute.go | 26 ++- .../execution/runtime_compute_resident.go | 220 ++++++++++++++++++ .../internal/execution/runtime_lifecycle.go | 6 +- .../internal/sandbox/e2b/provider.go | 15 ++ .../internal/sandbox/providers/registry.go | 1 + .../sandbox/providers/registry_test.go | 2 +- .../internal/sandbox/sandbox_provider.go | 10 + .../store/runtime_compute_lifecycle_test.go | 17 +- .../store/runtime_idle_policy_test.go | 12 + .../store/runtime_resident_lifecycle_test.go | 149 ++++++++++++ .../internal/store/runtime_suspension.go | 23 +- .../store/sandbox_deployment_mutations.go | 11 + .../store/sandbox_deployment_setup.go | 2 +- .../store/sandbox_deployment_view_test.go | 17 +- .../agents-api/tools/e2b-provider/README.md | 8 +- .../agents-api/tools/e2b-provider/provider.py | 50 +++- .../tools/e2b-provider/provider_test.py | 32 +++ 22 files changed, 627 insertions(+), 27 deletions(-) create mode 100644 services/agents-api/internal/execution/runtime_compute_resident.go create mode 100644 services/agents-api/internal/store/runtime_resident_lifecycle_test.go diff --git a/apps/docs/content/docs/sandbox-provider.mdx b/apps/docs/content/docs/sandbox-provider.mdx index 8e221c313..9e9f336a3 100644 --- a/apps/docs/content/docs/sandbox-provider.mdx +++ b/apps/docs/content/docs/sandbox-provider.mdx @@ -75,6 +75,7 @@ vendor-specific execution path. A backend without a native renewable lease | --- | --- | --- | | `sandbox.SandboxProvider` | Required | `Create`, `GetInfo`, `Renew`, `Kill`, and bounded bootstrap/diagnostic `RunCommand` | | `sandbox.CheckpointProvider` | Optional, separate interface | Exact compute incarnations, snapshot capture/restore, retained-source resume and cleanup | +| `sandbox.ResidentPauseProvider` | Optional, separate interface | Memory-preserving pause and resume of the same owned compute ID | | `runtimeobs.Source` | Optional, separate interface | Read-only, ownership-checked resource observations | | `runtimeobs.BatchSource` | Optional, separate interface | Bounded observations in input order, with per-target errors; `ok=false` means no batch read occurred | @@ -154,6 +155,12 @@ and snapshot, not whichever instance currently has the same display name. See [the lifecycle implementation](https://github.com/MiniMax-AI/parsar-core/blob/f6d258735fc601c521dd990e6f9e1ed261f4ef2d/services/agents-api/internal/execution/runtime_compute.go) and its failure tests before advertising this capability. +Resident pause retains the original provider ID and has no separate snapshot +identity. Core quiesces the Runtime before calling `Pause`, persists the phase +before the native call, and admits work only after `Resume` and Runtime +reconnection. An unknown pause result must be observed without replaying a +late mutation. See [the resident lifecycle](https://github.com/MiniMax-AI/parsar-core/blob/f6d258735fc601c521dd990e6f9e1ed261f4ef2d/services/agents-api/internal/execution/runtime_compute_resident.go). + ### Four distinct readiness facts | Fact | Evidence | Does not establish | @@ -230,8 +237,9 @@ proxy exposes checkpoint operations only for its registered checkpoint-capable backend. A new node backend must register both construction and the corresponding proxy capability; otherwise that capability is unavailable. Provider names belong in this adapter/configuration wiring, not Session scheduling, capability -preparation or Turn execution. Common lifecycle code uses `CheckpointProvider` -to admit suspension, independently of a provider name. +preparation or Turn execution. Common lifecycle code uses either +`CheckpointProvider` or `ResidentPauseProvider` to admit suspension, +independently of a provider name. The backend fingerprint identifies a native resource namespace, not mutable capacity. Core retains deployment generations so old owned allocations continue diff --git a/apps/docs/content/guide-sources.json b/apps/docs/content/guide-sources.json index 3bb814ba3..21a79edc2 100644 --- a/apps/docs/content/guide-sources.json +++ b/apps/docs/content/guide-sources.json @@ -28,7 +28,7 @@ "contracts/agents-api/harness-onboarding.md": "874a80dc1ffc5e97b2783bea3f6b217893f997b38bb8180a29f09c6dbf75e622", "docs/runtime-bootstrap.md": "0d49aed73b298039e04453e6f465b0e925fb2227fa35d4206820bb3acbd8df39", "docs/runtime-protocol.md": "e8aa4cf862b5f63a4138cfeb25196a6bdb5c40a2e389e828b464585415dc6c45", - "docs/sandbox-provider.md": "4d47b234a457f874f3ff7e61a7f5da6fc8b75b0c6df0ce0ce9ead1c07afdfb3e", + "docs/sandbox-provider.md": "e19f4947bb51edd9a3a9fefb2b0a9af91adcca822f1e9b6d1066c2426c575649", "apps/docs/scripts/guides.json": "3c3768fb94fd3d42464d4ce8c55724b9ba8360133c7cc19d513d8ec83a8e034f" }, "outputs": { @@ -59,6 +59,6 @@ "content/docs/harness-onboarding.mdx": "17626e9f6256ddfffc95fd81dfa8d9c712d9ff16792ea5ef54bfe8c3ce1a2a9f", "content/docs/runtime-bootstrap.mdx": "58982955811c9ad46a9fc04d8d2aef5a762fc9d25d5833f88aa75ee3fc46523f", "content/docs/runtime-protocol.mdx": "d6b818edf2a0e0c4f2d3878f75805043c6dba3b09c6563f80faa4cecd05d40c7", - "content/docs/sandbox-provider.mdx": "9db833fe9ff3f30a65451f80d7251348695a08e7830f9e95871739a710529818" + "content/docs/sandbox-provider.mdx": "71a68a37dbe443655cd7df472ff269b2bb0a065767b2e9dd32442fd57e894e2b" } } diff --git a/contracts/agents-api/sandbox-deployment.md b/contracts/agents-api/sandbox-deployment.md index 7bef7f43f..13c53dd78 100644 --- a/contracts/agents-api/sandbox-deployment.md +++ b/contracts/agents-api/sandbox-deployment.md @@ -207,9 +207,20 @@ enforce separately. Unknown values are null, including every value of a selectio saved before Core recorded them; a verified write records them. An omitted-key identical PUT is a no-op and does not refresh provider metadata. -`suspension` is `{idle_seconds, retention_seconds}` for microsandbox, the only -provider Core suspends (currently 300 and 86400). Docker, E2B and unconfigured -deployments return null. +`suspension` is `{idle_seconds, retention_seconds}` for microsandbox and new E2B +selections (currently 300 and 86400). Docker and unconfigured deployments return +null. E2B deployments saved before this policy retain `null` until an operator +submits the same E2B selection with an updated `expected_generation` through +`PUT /core/v1/sandbox/deployment`; that explicit update enables the policy and +increments the deployment generation. +After initialization, Core pauses an idle E2B sandbox with its memory retained, +including a Session that has not yet received a Turn. +The next input resumes the same sandbox ID and waits for the Runtime daemon to +reconnect before admission. The 86400-second retention bounds the paused +allocation; cleanup after expiry removes the sandbox, so a later request cannot +resume that original compute. Paused provider storage may still incur charges. +An uncertain pause result is observed under the allocation's private receipt, +never replayed against a potentially late native request. Two similarly named fields have different purposes: @@ -456,7 +467,7 @@ fields from the administrator node routes; runtime fields from the | Deployment `specification.resources` | `cpus` and `memory_mib`, equal to the ready template build's and taken from it when omitted; no disk fields | `cpus` and `memory_mib`; no disk quota | `cpus`, `memory_mib`, `root_disk_mib` and `environment_disk_mib` | | Deployment `specification.runtime` | Absent; the build is selected by `e2b.template` | The full [release](#runtime-release); nodes match `image_id` or `image_manifest_digest` | The full [release](#runtime-release); nodes match `microsandbox_ref`, `runtime_sha256` and `firmware_sha256` | | Deployment `e2b.template_build` | The build as Core read it when the selection was saved | Absent, with the whole `e2b` object | Absent, with the whole `e2b` object | -| Deployment `suspension` | `null`; Core does not suspend E2B sandboxes | `null` | `{idle_seconds, retention_seconds}` | +| Deployment `suspension` | `{idle_seconds, retention_seconds}` | `null` | `{idle_seconds, retention_seconds}` | | Deployment `resources.allocations`, `resources.pending` | Core's unreleased E2B sandboxes, and hosted Environments waiting for one | Totals across all nodes | Totals across all nodes | | Enrollment-token `max_active`, `max_retained` | 409 `sandbox_deployment_conflict`, after the 400 capacity checks; E2B has no nodes | `max_retained` always equals `max_active` | Both limits apply | | Node list and detail | Empty list; detail returns 404 | Enrolled nodes | Enrolled nodes | @@ -466,7 +477,7 @@ fields from the administrator node routes; runtime fields from the | Allocation `compute_phase`, `compute_phase_changed_at` | Not applicable: no node allocations | Always `disabled`, counted as running until release; the time is the allocation's creation | Includes `suspended`; its time plus `suspension.retention_seconds` tells roughly when Core reclaims the snapshot | | Runtime observation `cpu`, `memory` | From E2B metrics: `cpu.utilization_ratio` and `capacity_cores`, memory usage and limit; no cumulative CPU time | From Docker stats: `cpu.usage_seconds_total`, CPU and memory limits, memory usage | From the VM: `cpu.usage_seconds_total`, CPU and memory limits, memory usage | | Runtime observation `disk` | E2B `diskUsed` and `diskTotal`; `null` when the template does not report them | `null`: no disk quota | `null` for now | -| Runtime observation `lifecycle_state: sleeping` | Never | Never | While suspended | +| Runtime observation `lifecycle_state: sleeping` | Not yet reported by E2B metrics; paused compute has no running sample | Never | While suspended | | Runtime history CPU | Mean of the utilization ratios E2B reported in each bucket | Derived from cumulative CPU time | Derived from cumulative CPU time | ## Failure and version boundaries diff --git a/docs/sandbox-provider.md b/docs/sandbox-provider.md index 60a6af2d7..0519010f6 100644 --- a/docs/sandbox-provider.md +++ b/docs/sandbox-provider.md @@ -72,6 +72,7 @@ vendor-specific execution path. A backend without a native renewable lease | --- | --- | --- | | `sandbox.SandboxProvider` | Required | `Create`, `GetInfo`, `Renew`, `Kill`, and bounded bootstrap/diagnostic `RunCommand` | | `sandbox.CheckpointProvider` | Optional, separate interface | Exact compute incarnations, snapshot capture/restore, retained-source resume and cleanup | +| `sandbox.ResidentPauseProvider` | Optional, separate interface | Memory-preserving pause and resume of the same owned compute ID | | `runtimeobs.Source` | Optional, separate interface | Read-only, ownership-checked resource observations | | `runtimeobs.BatchSource` | Optional, separate interface | Bounded observations in input order, with per-target errors; `ok=false` means no batch read occurred | @@ -151,6 +152,12 @@ and snapshot, not whichever instance currently has the same display name. See [the lifecycle implementation](../services/agents-api/internal/execution/runtime_compute.go) and its failure tests before advertising this capability. +Resident pause retains the original provider ID and has no separate snapshot +identity. Core quiesces the Runtime before calling `Pause`, persists the phase +before the native call, and admits work only after `Resume` and Runtime +reconnection. An unknown pause result must be observed without replaying a +late mutation. See [the resident lifecycle](../services/agents-api/internal/execution/runtime_compute_resident.go). + ### Four distinct readiness facts | Fact | Evidence | Does not establish | @@ -227,8 +234,9 @@ proxy exposes checkpoint operations only for its registered checkpoint-capable backend. A new node backend must register both construction and the corresponding proxy capability; otherwise that capability is unavailable. Provider names belong in this adapter/configuration wiring, not Session scheduling, capability -preparation or Turn execution. Common lifecycle code uses `CheckpointProvider` -to admit suspension, independently of a provider name. +preparation or Turn execution. Common lifecycle code uses either +`CheckpointProvider` or `ResidentPauseProvider` to admit suspension, +independently of a provider name. The backend fingerprint identifies a native resource namespace, not mutable capacity. Core retains deployment generations so old owned allocations continue diff --git a/services/agents-api/cmd/server/managed_setup.go b/services/agents-api/cmd/server/managed_setup.go index 2ab0c49f6..8c125e172 100644 --- a/services/agents-api/cmd/server/managed_setup.go +++ b/services/agents-api/cmd/server/managed_setup.go @@ -121,7 +121,9 @@ func (s *managedSetup) configuration(setup store.SandboxSetup) (execution.Prepar } selected := &execution.RuntimeProvider{InstallationID: setup.InstallationID, ProviderKind: setup.Provider, Generation: setup.Generation, Mode: setup.Mode, AdmissionPaused: setup.AdmissionPaused, CoreURL: s.publicURL + "/api/v1", BackendFingerprint: setup.BackendFingerprint, Provider: provider} - if _, supportsCheckpoint := provider.(sandbox.CheckpointProvider); supportsCheckpoint { + _, supportsCheckpoint := provider.(sandbox.CheckpointProvider) + _, supportsResidentPause := provider.(sandbox.ResidentPauseProvider) + if (supportsCheckpoint || supportsResidentPause) && setup.IdleSeconds > 0 && setup.RetentionSeconds > 0 { selected.Suspension = &execution.RuntimeSuspensionPolicy{IdleTimeout: time.Duration(setup.IdleSeconds) * time.Second, Retention: time.Duration(setup.RetentionSeconds) * time.Second, MaxActive: 4, MaxRetained: 16} } diff --git a/services/agents-api/internal/execution/runtime_compute.go b/services/agents-api/internal/execution/runtime_compute.go index 7fc56384a..2f3cdd4fb 100644 --- a/services/agents-api/internal/execution/runtime_compute.go +++ b/services/agents-api/internal/execution/runtime_compute.go @@ -12,8 +12,8 @@ import ( "github.com/google/uuid" ) -// RuntimeSuspensionPolicy applies only to an explicitly qualified single-host -// provider. Fixed guest sizing plus MaxActive bounds reserved CPU and memory. +// RuntimeSuspensionPolicy applies only to a qualified pause-capable provider. +// Node-backed checkpoint providers also use MaxActive to bound reserved capacity. type RuntimeSuspensionPolicy struct { IdleTimeout time.Duration Retention time.Duration @@ -35,7 +35,7 @@ func (r *runtimeLifecycle) computeCapacity(ctx context.Context, key string) erro if key != r.config.InstallationID { return sandbox.ErrOwnership } - if policy == nil { + if policy == nil || r.config.Mode == "direct" { return nil } count, err := r.store.CountRuntimeComputeReservations(ctx, key) @@ -60,9 +60,26 @@ func (r *runtimeLifecycle) saveCompute(ctx context.Context, owner store.RuntimeA } idleTimeout = policy.IdleTimeout } + _, resident := r.config.Provider.(sandbox.ResidentPauseProvider) + if resident { + return r.store.SetRuntimeResidentCompute(ctx, owner, phase, raw, until, idleTimeout) + } return r.store.SetRuntimeCompute(ctx, owner, phase, raw, until, idleTimeout) } func (r *runtimeLifecycle) enableCompute(ctx context.Context, owner store.RuntimeAllocation) error { + if p, ok := r.config.Provider.(sandbox.ResidentPauseProvider); ok { + info, err := p.GetInfo(ctx, runtimeReference(owner)) + if err != nil { + return err + } + if info.Reference != runtimeReference(owner) || info.State != "running" || !info.BootstrapComplete || info.ProviderID == "" { + return sandbox.ErrComputeUnconfirmed + } + if _, err = r.saveCompute(ctx, owner, "running", runtimeCompute{Current: sandbox.Compute{ID: info.ProviderID}}, nil); err != nil { + return err + } + return r.store.TouchRuntimeActivity(ctx, owner.TenantID, owner.EnvironmentID) + } p, ok := r.config.Provider.(sandbox.CheckpointProvider) if !ok { return sandbox.ErrInvalid @@ -83,6 +100,9 @@ func (r *runtimeLifecycle) enableCompute(ctx context.Context, owner store.Runtim } func (r *runtimeLifecycle) observeCompute(ctx context.Context, owner store.RuntimeAllocation) error { + if p, ok := r.config.Provider.(sandbox.ResidentPauseProvider); ok { + return r.observeResidentCompute(ctx, p, owner) + } p, ok := r.config.Provider.(sandbox.CheckpointProvider) if !ok { return sandbox.ErrInvalid diff --git a/services/agents-api/internal/execution/runtime_compute_resident.go b/services/agents-api/internal/execution/runtime_compute_resident.go new file mode 100644 index 000000000..a987f1cc6 --- /dev/null +++ b/services/agents-api/internal/execution/runtime_compute_resident.go @@ -0,0 +1,220 @@ +package execution + +import ( + "context" + "encoding/json" + "errors" + "time" + + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/gateway" + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/proto" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" + "github.com/google/uuid" +) + +// Resident suspension retains the original provider ID. It intentionally does +// not claim a checkpoint or use the snapshot path, which kills its source VM. +func (r *runtimeLifecycle) observeResidentCompute(ctx context.Context, p sandbox.ResidentPauseProvider, owner store.RuntimeAllocation) error { + var state runtimeCompute + if !validResidentState(owner, &state) { + return sandbox.ErrOwnership + } + if owner.SessionDeleted || owner.Expired || owner.State == "cleanup_pending" { + if err := r.store.CheckExecutionOwnership(ctx); err != nil { + return err + } + if err := p.Kill(ctx, runtimeReference(owner)); err != nil { + return err + } + _, err := r.store.ReleaseRuntimeAllocation(ctx, owner) + return err + } + if err := r.store.CheckExecutionOwnership(ctx); err != nil { + return err + } + switch owner.ComputePhase { + case "running": + return r.idleResidentCompute(ctx, p, owner, state) + case "quiescing": + state.Rollback = true + next, err := r.saveCompute(ctx, owner, "waking", state, owner.ComputeRetainedUntil) + if err != nil { + return err + } + return r.wakeResidentCompute(ctx, p, next, state) + case "suspending": + info, err := p.Pause(ctx, runtimeReference(owner)) + if err != nil { + return err + } + if !residentInfo(owner, state, info, "paused") { + return sandbox.ErrComputeUnconfirmed + } + _, err = r.saveCompute(ctx, owner, "suspended", state, owner.ComputeRetainedUntil) + return err + case "suspended": + activity, err := r.store.RuntimeActivity(ctx, owner) + if err != nil || !activity.Busy && !activity.WakeRequested { + return err + } + if err := r.computeCapacityForAllocation(ctx, owner); err != nil { + return err + } + state.RestoreID = uuid.NewString() + next, err := r.saveCompute(ctx, owner, "restoring", state, owner.ComputeRetainedUntil) + if err != nil { + return err + } + return r.restoreResidentCompute(ctx, p, next, state) + case "restoring": + return r.restoreResidentCompute(ctx, p, owner, state) + case "waking": + return r.wakeResidentCompute(ctx, p, owner, state) + default: + return sandbox.ErrInvalid + } +} + +func validResidentState(owner store.RuntimeAllocation, state *runtimeCompute) bool { + return len(owner.ComputeState) > 0 && json.Unmarshal(owner.ComputeState, state) == nil && + state.Current.ID != "" && state.Target == nil && state.Snapshot == nil +} + +func residentInfo(owner store.RuntimeAllocation, state runtimeCompute, info sandbox.Info, status string) bool { + return info.Reference == runtimeReference(owner) && info.ProviderID == state.Current.ID && + info.State == status && info.BootstrapComplete && info.CreateSettled +} + +func (r *runtimeLifecycle) idleResidentCompute(ctx context.Context, p sandbox.ResidentPauseProvider, owner store.RuntimeAllocation, state runtimeCompute) error { + activity, err := r.store.RuntimeActivity(ctx, owner) + if err != nil { + return err + } + if activity.WakeRequested { + return r.store.ClearRuntimeWake(ctx, owner, owner.ComputeActivityAt) + } + policy := r.config.Suspension + if policy == nil || !activity.ReadyToPauseResident(policy.IdleTimeout) { + info, err := p.Renew(ctx, runtimeReference(owner)) + if err != nil { + return err + } + if !residentInfo(owner, state, info, "running") { + return sandbox.ErrComputeUnconfirmed + } + _, err = r.store.KeepRuntimeAllocation(ctx, owner) + return err + } + info, err := p.GetInfo(ctx, runtimeReference(owner)) + if err != nil { + return err + } + if !residentInfo(owner, state, info, "running") { + return sandbox.ErrComputeUnconfirmed + } + peer, err := authorizedRuntimePeer(ctx, r.store, r.registry, owner.DeviceID) + if err != nil { + return err + } + state.SuspendID, state.RestoreID, state.Rollback = uuid.NewString(), "", false + until := activity.ObservedAt.Add(policy.Retention) + next, err := r.saveCompute(ctx, owner, "quiescing", state, &until) + if err != nil { + return err + } + result, err := peer.SuspendControl(ctx, proto.TypeEnvironmentQuiesce, proto.EnvironmentSuspendPayload{EnvironmentID: owner.EnvironmentID, SuspendID: state.SuspendID}) + if err != nil { + return err + } + if !result.Accepted { + if err := r.store.TouchRuntimeActivity(ctx, owner.TenantID, owner.EnvironmentID); err != nil { + return err + } + _, err = r.saveCompute(ctx, next, "running", runtimeCompute{Current: state.Current}, nil) + return err + } + if err := observeRuntimeConnection(ctx, r.store, r.connections, owner.TenantID, owner.EnvironmentID, nil, false); err != nil { + return err + } + suspending, err := r.saveCompute(ctx, next, "suspending", state, &until) + if errors.Is(err, store.ErrTurnConflict) { + state.Rollback = true + next, err = r.saveCompute(ctx, next, "waking", state, &until) + if err != nil { + return err + } + return r.wakeResidentCompute(ctx, p, next, state) + } + if err != nil { + return err + } + info, err = p.Pause(ctx, runtimeReference(owner)) + if err != nil { + return err + } + if !residentInfo(owner, state, info, "paused") { + return sandbox.ErrComputeUnconfirmed + } + _, err = r.saveCompute(ctx, suspending, "suspended", state, &until) + return err +} + +func (r *runtimeLifecycle) restoreResidentCompute(ctx context.Context, p sandbox.ResidentPauseProvider, owner store.RuntimeAllocation, state runtimeCompute) error { + info, err := p.Resume(ctx, runtimeReference(owner)) + if err != nil { + return err + } + if !residentInfo(owner, state, info, "running") { + return sandbox.ErrComputeUnconfirmed + } + next, err := r.saveCompute(ctx, owner, "waking", state, owner.ComputeRetainedUntil) + if err != nil { + return err + } + return r.wakeResidentCompute(ctx, p, next, state) +} + +func (r *runtimeLifecycle) wakeResidentCompute(ctx context.Context, p sandbox.ResidentPauseProvider, owner store.RuntimeAllocation, state runtimeCompute) error { + peer, err := authorizedRuntimePeer(ctx, r.store, r.registry, owner.DeviceID) + if err != nil { + if !errors.Is(err, store.ErrNotFound) && !errors.Is(err, gateway.ErrDeviceNotRegistered) && !errors.Is(err, gateway.ErrSessionClosed) { + return err + } + result, err := p.RunCommand(ctx, runtimeReference(owner), sandbox.Command{Args: []string{"oac-daemon", "resume", "--control-file", "/run/oac/daemon-suspend.json", "--environment-id", owner.EnvironmentID, "--suspend-id", state.SuspendID}}) + if err != nil { + return err + } + if result.ExitCode != 0 { + return sandbox.ErrComputeUnconfirmed + } + timer := time.NewTicker(100 * time.Millisecond) + defer timer.Stop() + for { + peer, err = authorizedRuntimePeer(ctx, r.store, r.registry, owner.DeviceID) + if err == nil { + break + } + select { + case <-ctx.Done(): + return ctx.Err() + case <-timer.C: + } + } + } + result, err := peer.SuspendControl(ctx, proto.TypeEnvironmentResume, proto.EnvironmentSuspendPayload{EnvironmentID: owner.EnvironmentID, SuspendID: state.SuspendID, Rollback: state.Rollback}) + if err != nil { + return err + } + if !result.Accepted { + return sandbox.ErrComputeUnconfirmed + } + next, err := r.saveCompute(ctx, owner, "running", runtimeCompute{Current: state.Current}, nil) + if err != nil { + return err + } + if err := r.store.ClearRuntimeWake(ctx, next, owner.ComputeActivityAt); err != nil { + return err + } + return r.observeConnection(ctx, next) +} diff --git a/services/agents-api/internal/execution/runtime_lifecycle.go b/services/agents-api/internal/execution/runtime_lifecycle.go index d372a4df9..8846a9bd3 100644 --- a/services/agents-api/internal/execution/runtime_lifecycle.go +++ b/services/agents-api/internal/execution/runtime_lifecycle.go @@ -93,12 +93,14 @@ func validatedRuntimeProvider(config *RuntimeProvider, registry *gateway.Registr if copied.Mode != "" && copied.Mode != "nodes" && copied.Mode != "direct" { return RuntimeProvider{}, sandbox.ErrInvalid } - if copied.Mode == "direct" && (copied.ProviderKind == "" || copied.LocalNodeID != "" || copied.Suspension != nil) { + if copied.Mode == "direct" && (copied.ProviderKind == "" || copied.LocalNodeID != "") { return RuntimeProvider{}, sandbox.ErrInvalid } if config.Suspension != nil { policy := *config.Suspension - if _, ok := config.Provider.(sandbox.CheckpointProvider); !ok || policy.IdleTimeout < time.Second || policy.Retention < time.Second || policy.MaxActive < 1 || policy.MaxRetained < policy.MaxActive { + _, checkpoint := config.Provider.(sandbox.CheckpointProvider) + _, resident := config.Provider.(sandbox.ResidentPauseProvider) + if (!checkpoint && !resident) || (resident && copied.Mode != "direct") || policy.IdleTimeout < time.Second || policy.Retention < time.Second || policy.MaxActive < 1 || policy.MaxRetained < policy.MaxActive { return RuntimeProvider{}, sandbox.ErrInvalid } copied.Suspension = &policy diff --git a/services/agents-api/internal/sandbox/e2b/provider.go b/services/agents-api/internal/sandbox/e2b/provider.go index 8d50db340..932f0c1aa 100644 --- a/services/agents-api/internal/sandbox/e2b/provider.go +++ b/services/agents-api/internal/sandbox/e2b/provider.go @@ -148,6 +148,7 @@ type Provider struct { } var _ sandbox.SandboxProvider = (*Provider)(nil) +var _ sandbox.ResidentPauseProvider = (*Provider)(nil) func validID(value string) bool { id, err := uuid.Parse(value) @@ -292,6 +293,20 @@ func (p *Provider) GetInfo(ctx context.Context, r sandbox.Reference) (sandbox.In func (p *Provider) Renew(ctx context.Context, r sandbox.Reference) (sandbox.Info, error) { return p.info(ctx, "renew", r, nil) } +func (p *Provider) Pause(ctx context.Context, r sandbox.Reference) (sandbox.Info, error) { + info, err := p.info(ctx, "pause", r, nil) + if err == nil && (info.Reference != r || info.ProviderID == "" || info.State != "paused" || !info.BootstrapComplete) { + return info, sandbox.ErrComputeUnconfirmed + } + return info, err +} +func (p *Provider) Resume(ctx context.Context, r sandbox.Reference) (sandbox.Info, error) { + info, err := p.info(ctx, "resume", r, nil) + if err == nil && (info.Reference != r || info.ProviderID == "" || info.State != "running" || !info.BootstrapComplete) { + return info, sandbox.ErrComputeUnconfirmed + } + return info, err +} func (p *Provider) Kill(ctx context.Context, r sandbox.Reference) error { out, err := p.call(ctx, "kill", r, nil, nil) if err == nil && (out.Info == nil || !out.Info.CreateSettled || out.Info.State != "absent") { diff --git a/services/agents-api/internal/sandbox/providers/registry.go b/services/agents-api/internal/sandbox/providers/registry.go index a06a2c30e..86c918ed6 100644 --- a/services/agents-api/internal/sandbox/providers/registry.go +++ b/services/agents-api/internal/sandbox/providers/registry.go @@ -48,6 +48,7 @@ var adapters = map[string]Adapter{ }, "e2b": { Policy: e2b.Policy(), Mode: "direct", BuildDirect: buildE2B, + IdleSeconds: 300, RetentionSeconds: 86400, Credential: true, PublicOrigin: true, ReplaceCredential: e2b.ReplaceCredential, CredentialRequiresReset: e2b.CredentialRequiresReset, CredentialUnconfirmed: e2b.ErrRequestUnconfirmed, Restore: e2b.RestoreSelection, diff --git a/services/agents-api/internal/sandbox/providers/registry_test.go b/services/agents-api/internal/sandbox/providers/registry_test.go index 17212d63a..e18aaa51c 100644 --- a/services/agents-api/internal/sandbox/providers/registry_test.go +++ b/services/agents-api/internal/sandbox/providers/registry_test.go @@ -17,7 +17,7 @@ func TestRegistrationOwnsDeploymentPolicy(t *testing.T) { }{ {"docker", "nodes", "nodes", 0, 0, false}, {"microsandbox", "nodes", "nodes", 300, 86400, true}, - {"e2b", "direct", "e2b", 0, 0, false}, + {"e2b", "direct", "e2b", 300, 86400, false}, } { t.Run(tc.kind, func(t *testing.T) { d, err := Describe(tc.kind, installation) diff --git a/services/agents-api/internal/sandbox/sandbox_provider.go b/services/agents-api/internal/sandbox/sandbox_provider.go index 531fd3d71..ad08747a0 100644 --- a/services/agents-api/internal/sandbox/sandbox_provider.go +++ b/services/agents-api/internal/sandbox/sandbox_provider.go @@ -89,6 +89,16 @@ type SandboxProvider interface { RunCommand(context.Context, Reference, Command) (CommandResult, error) } +// ResidentPauseProvider keeps the same owned compute incarnation across a +// memory-preserving pause. Pause must never create compute, and Resume must +// reconnect only the existing allocation. Both operations must reconcile an +// unknown result against provider state before reporting success. +type ResidentPauseProvider interface { + SandboxProvider + Pause(context.Context, Reference) (Info, error) + Resume(context.Context, Reference) (Info, error) +} + // CheckpointProvider is optional. It supplements the existing provider with // exact-incarnation operations; Worker and Store remain the lifecycle owner. type CheckpointProvider interface { diff --git a/services/agents-api/internal/store/runtime_compute_lifecycle_test.go b/services/agents-api/internal/store/runtime_compute_lifecycle_test.go index e2665e687..bdc98a255 100644 --- a/services/agents-api/internal/store/runtime_compute_lifecycle_test.go +++ b/services/agents-api/internal/store/runtime_compute_lifecycle_test.go @@ -260,6 +260,7 @@ type computeLifecycleFixture struct { store *store.Store pool *pgxpool.Pool provider *fakeCheckpointProvider + managed sandbox.SandboxProvider worker *execution.Worker stop func() key string @@ -267,6 +268,10 @@ type computeLifecycleFixture struct { } func newComputeLifecycleFixture(t *testing.T, maxActive, maxRetained int) *computeLifecycleFixture { + return newComputeLifecycleFixtureWithProvider(t, maxActive, maxRetained, nil) +} + +func newComputeLifecycleFixtureWithProvider(t *testing.T, maxActive, maxRetained int, wrap func(*fakeCheckpointProvider) sandbox.SandboxProvider) *computeLifecycleFixture { t.Helper() s, pool := store.NewManagedTestStore(t) registry := gateway.NewRegistry() @@ -282,14 +287,22 @@ func newComputeLifecycleFixture(t *testing.T, maxActive, maxRetained int) *compu } server.Close() }) - f := &computeLifecycleFixture{t: t, store: s, pool: pool, provider: p, key: uuid.NewString(), policy: execution.RuntimeSuspensionPolicy{IdleTimeout: time.Second, Retention: time.Hour, MaxActive: maxActive, MaxRetained: maxRetained}} + managed := sandbox.SandboxProvider(p) + if wrap != nil { + managed = wrap(p) + } + f := &computeLifecycleFixture{t: t, store: s, pool: pool, provider: p, managed: managed, key: uuid.NewString(), policy: execution.RuntimeSuspensionPolicy{IdleTimeout: time.Second, Retention: time.Hour, MaxActive: maxActive, MaxRetained: maxRetained}} f.start() return f } func (f *computeLifecycleFixture) start() { t := f.t t.Helper() - w, err := execution.StartWorker(t.Context(), &execution.Dispatcher{Store: f.store, Registry: f.provider.registry, ManagedRuntimes: &execution.RuntimeProvider{CoreURL: "http://core.invalid/api/v1", InstallationID: f.key, BackendFingerprint: "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", Provider: f.provider, Suspension: &f.policy}}) + config := &execution.RuntimeProvider{CoreURL: "http://core.invalid/api/v1", InstallationID: f.key, BackendFingerprint: "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", Provider: f.managed, Suspension: &f.policy} + if _, resident := f.managed.(sandbox.ResidentPauseProvider); resident { + config.Mode, config.ProviderKind = "direct", "e2b" + } + w, err := execution.StartWorker(t.Context(), &execution.Dispatcher{Store: f.store, Registry: f.provider.registry, ManagedRuntimes: config}) if err != nil { t.Fatal(err) } diff --git a/services/agents-api/internal/store/runtime_idle_policy_test.go b/services/agents-api/internal/store/runtime_idle_policy_test.go index b08b740cf..2922240e6 100644 --- a/services/agents-api/internal/store/runtime_idle_policy_test.go +++ b/services/agents-api/internal/store/runtime_idle_policy_test.go @@ -37,3 +37,15 @@ func TestRuntimeIdleAdmissionUsesDatabaseClock(t *testing.T) { }) } } + +func TestResidentIdleAdmissionIncludesNeverUsedSessions(t *testing.T) { + now := time.Now() + activity := RuntimeActivity{ObservedAt: now, LastActivity: now.Add(-6 * time.Minute)} + if activity.ReadyToSuspend(5*time.Minute) || !activity.ReadyToPauseResident(5*time.Minute) { + t.Fatal("resident policy did not distinguish never-used Session") + } + activity.Busy = true + if activity.ReadyToPauseResident(5 * time.Minute) { + t.Fatal("resident policy paused pending work") + } +} diff --git a/services/agents-api/internal/store/runtime_resident_lifecycle_test.go b/services/agents-api/internal/store/runtime_resident_lifecycle_test.go new file mode 100644 index 000000000..e76c77725 --- /dev/null +++ b/services/agents-api/internal/store/runtime_resident_lifecycle_test.go @@ -0,0 +1,149 @@ +package store_test + +import ( + "context" + "encoding/json" + "errors" + "testing" + + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" +) + +type fakeResidentProvider struct { + *fakeCheckpointProvider + pauses, resumes int + losePause bool +} + +func (p *fakeResidentProvider) Create(ctx context.Context, b sandbox.Bootstrap) (sandbox.Info, error) { + info, err := p.fakeCheckpointProvider.Create(ctx, b) + if err != nil { + return info, err + } + p.mu.Lock() + info.CreateSettled = true + p.resources[b.AllocationID] = info + p.mu.Unlock() + return info, nil +} +func (p *fakeResidentProvider) Pause(ctx context.Context, r sandbox.Reference) (sandbox.Info, error) { + info, err := p.GetInfo(ctx, r) + if err != nil { + return info, err + } + p.mu.Lock() + defer p.mu.Unlock() + if info.State != "running" { + return info, errors.New("resident pause replayed") + } + p.pauses++ + info.State = "paused" + p.resources[r.AllocationID] = info + if p.losePause { + p.losePause = false + return sandbox.Info{}, sandbox.ErrComputeUnconfirmed + } + return info, nil +} +func (p *fakeResidentProvider) Resume(ctx context.Context, r sandbox.Reference) (sandbox.Info, error) { + info, err := p.GetInfo(ctx, r) + if err != nil { + return info, err + } + p.mu.Lock() + defer p.mu.Unlock() + if info.State != "paused" { + return info, errors.New("resident resume replayed") + } + p.resumes++ + info.State = "running" + p.resources[r.AllocationID] = info + return info, nil +} +func (p *fakeResidentProvider) RunCommand(ctx context.Context, r sandbox.Reference, command sandbox.Command) (sandbox.CommandResult, error) { + if len(command.Args) != 8 || command.Args[0] != "oac-daemon" || command.Args[1] != "resume" || command.Args[5] != r.EnvironmentID { + return sandbox.CommandResult{}, errors.New("unexpected resident wake command") + } + p.mu.Lock() + p.wakeCommands++ + b := p.bootstraps[r.AllocationID] + p.mu.Unlock() + return sandbox.CommandResult{}, p.connect(ctx, b) +} + +func TestResidentIdlePauseAndWakeKeepOriginalCompute(t *testing.T) { + var resident *fakeResidentProvider + f := newComputeLifecycleFixtureWithProvider(t, 2, 4, func(p *fakeCheckpointProvider) sandbox.SandboxProvider { + resident = &fakeResidentProvider{fakeCheckpointProvider: p} + return resident + }) + tenant, _, environment, owner := f.create() + var before struct{ Current sandbox.Compute } + if err := json.Unmarshal(owner.ComputeState, &before); err != nil { + t.Fatal(err) + } + f.complete(owner) + owner = f.phase(tenant, environment.ID, "suspended") + if resident.pauses != 1 || resident.resumes != 0 || f.provider.computeKills != 0 || f.provider.captures != 0 { + t.Fatal("resident pause created or destroyed compute", resident.pauses, f.provider.computeKills) + } + f.queued(owner) + owner = f.phase(tenant, environment.ID, "running") + var after struct{ Current sandbox.Compute } + if err := json.Unmarshal(owner.ComputeState, &after); err != nil { + t.Fatal(err) + } + if before.Current.ID == "" || before.Current.ID != after.Current.ID || resident.resumes != 1 || f.provider.wakeCommands != 1 { + t.Fatal("original compute or daemon connection was not restored") + } +} + +func TestResidentUnknownPauseUsesObservationWithoutReplay(t *testing.T) { + var resident *fakeResidentProvider + f := newComputeLifecycleFixtureWithProvider(t, 2, 4, func(p *fakeCheckpointProvider) sandbox.SandboxProvider { + resident = &fakeResidentProvider{fakeCheckpointProvider: p, losePause: true} + return resident + }) + tenant, _, environment, owner := f.create() + f.phase(tenant, environment.ID, "running") + f.complete(owner) + if err := f.worker.ReconcileManagedRuntimes(t.Context()); !errors.Is(err, sandbox.ErrComputeUnconfirmed) { + t.Fatal("unknown pause was reported as settled", err) + } + f.phase(tenant, environment.ID, "suspended") + if resident.pauses != 1 { + t.Fatal("unknown pause was replayed", resident.pauses) + } +} + +func TestResidentDoesNotPauseWhileTurnQueued(t *testing.T) { + var resident *fakeResidentProvider + f := newComputeLifecycleFixtureWithProvider(t, 2, 4, func(p *fakeCheckpointProvider) sandbox.SandboxProvider { + resident = &fakeResidentProvider{fakeCheckpointProvider: p} + return resident + }) + tenant, _, environment, owner := f.create() + f.complete(owner) + f.queued(owner) + f.phase(tenant, environment.ID, "running") + if resident.pauses != 0 { + t.Fatal("busy session was paused") + } +} + +func TestResidentNeverUsedSessionPausesAfterInitializationIdle(t *testing.T) { + var resident *fakeResidentProvider + f := newComputeLifecycleFixtureWithProvider(t, 2, 4, func(p *fakeCheckpointProvider) sandbox.SandboxProvider { + resident = &fakeResidentProvider{fakeCheckpointProvider: p} + return resident + }) + tenant, _, environment, owner := f.create() + // The initial Runtime-ready touch prevents an immediate pause even when + // provisioning took longer than the idle threshold. + f.phase(tenant, environment.ID, "running") + f.sql(`UPDATE runtime_allocations SET compute_activity_at=clock_timestamp()-interval '10 minutes',created_at=clock_timestamp()-interval '10 minutes' WHERE id=$1`, owner.ID) + f.phase(tenant, environment.ID, "suspended") + if resident.pauses != 1 { + t.Fatal("never-used idle Session kept a running sandbox") + } +} diff --git a/services/agents-api/internal/store/runtime_suspension.go b/services/agents-api/internal/store/runtime_suspension.go index 67b922d89..5921f3983 100644 --- a/services/agents-api/internal/store/runtime_suspension.go +++ b/services/agents-api/internal/store/runtime_suspension.go @@ -23,6 +23,13 @@ func (a RuntimeActivity) ReadyToSuspend(idleTimeout time.Duration) bool { a.ObservedAt.Sub(a.LastActivity) >= idleTimeout } +// A resident provider may also pause a Session that was initialized but never +// received a Turn. Its idle clock starts when the Runtime first becomes ready. +func (a RuntimeActivity) ReadyToPauseResident(idleTimeout time.Duration) bool { + return idleTimeout > 0 && !a.Busy && !a.WakeRequested && + a.ObservedAt.Sub(a.LastActivity) >= idleTimeout +} + func runtimeActivity(row sqlc.GetRuntimeActivityRow) RuntimeActivity { return RuntimeActivity{LastActivity: row.LastActivity.Time, ObservedAt: row.ObservedAt.Time, Busy: row.Busy, WakeRequested: row.ComputeWakeRequested, HasCompletedTurn: row.HasCompletedTurn} } @@ -30,6 +37,16 @@ func runtimeActivity(row sqlc.GetRuntimeActivityRow) RuntimeActivity { // SetRuntimeCompute commits an operation phase before its external effects. // Revision and the existing Session lock fence a stale lifecycle observation. func (s *Store) SetRuntimeCompute(ctx context.Context, owner RuntimeAllocation, phase string, state json.RawMessage, retainedUntil *time.Time, idleTimeout time.Duration) (RuntimeAllocation, error) { + return s.setRuntimeCompute(ctx, owner, phase, state, retainedUntil, idleTimeout, false) +} + +// SetRuntimeResidentCompute uses the same Session lock and phase checks while +// allowing an initialized Session without a completed Turn to become idle. +func (s *Store) SetRuntimeResidentCompute(ctx context.Context, owner RuntimeAllocation, phase string, state json.RawMessage, retainedUntil *time.Time, idleTimeout time.Duration) (RuntimeAllocation, error) { + return s.setRuntimeCompute(ctx, owner, phase, state, retainedUntil, idleTimeout, true) +} + +func (s *Store) setRuntimeCompute(ctx context.Context, owner RuntimeAllocation, phase string, state json.RawMessage, retainedUntil *time.Time, idleTimeout time.Duration, allowUnstarted bool) (RuntimeAllocation, error) { if !runtimeComputeTransition(owner.ComputePhase, phase) || !json.Valid(state) || (owner.ComputePhase == "running" && phase == "quiescing" && idleTimeout <= 0) { return RuntimeAllocation{}, ErrInvalidInput } @@ -52,7 +69,11 @@ func (s *Store) SetRuntimeCompute(ctx context.Context, owner RuntimeAllocation, if activity.Busy || activity.ComputeWakeRequested { return sqlc.RuntimeAllocation{}, ErrTurnConflict } - if phase == "quiescing" && (!runtimeActivity(activity).ReadyToSuspend(idleTimeout) || row.ComputeActivityAt.Time.After(owner.ComputeActivityAt)) { + ready := runtimeActivity(activity).ReadyToSuspend(idleTimeout) + if allowUnstarted { + ready = runtimeActivity(activity).ReadyToPauseResident(idleTimeout) + } + if phase == "quiescing" && (!ready || row.ComputeActivityAt.Time.After(owner.ComputeActivityAt)) { return sqlc.RuntimeAllocation{}, ErrTurnConflict } } diff --git a/services/agents-api/internal/store/sandbox_deployment_mutations.go b/services/agents-api/internal/store/sandbox_deployment_mutations.go index df0de2d55..8ff1cc576 100644 --- a/services/agents-api/internal/store/sandbox_deployment_mutations.go +++ b/services/agents-api/internal/store/sandbox_deployment_mutations.go @@ -39,6 +39,17 @@ func (s *Store) sandboxSelectionEqual(d sqlc.RuntimeDeployment, input SandboxDep if d.ProviderKind != input.Provider { return false, nil } + if input.Provider == "e2b" { + policy, err := providers.Describe(input.Provider, runtimeUUID(d.InstallationID)) + if err != nil { + return false, err + } + if d.IdleSeconds != policy.IdleSeconds || d.RetentionSeconds != policy.RetentionSeconds { + // A same-selection admin PUT explicitly opts an older deployment in + // without changing every installed E2B deployment during migration. + return false, nil + } + } if input.E2B == nil { return d.E2bTemplate == "", nil } diff --git a/services/agents-api/internal/store/sandbox_deployment_setup.go b/services/agents-api/internal/store/sandbox_deployment_setup.go index 40796a341..b4d675b35 100644 --- a/services/agents-api/internal/store/sandbox_deployment_setup.go +++ b/services/agents-api/internal/store/sandbox_deployment_setup.go @@ -180,7 +180,7 @@ func runtimeDeploymentView(d sqlc.RuntimeDeployment, publicURL string) RuntimeDe result.E2B.TemplateBuild.Status = &status } } - if providers.SupportsCheckpoint(d.ProviderKind) { + if d.IdleSeconds > 0 && d.RetentionSeconds > 0 { result.Suspension = &SandboxSuspensionView{IdleSeconds: d.IdleSeconds, RetentionSeconds: d.RetentionSeconds} } return result diff --git a/services/agents-api/internal/store/sandbox_deployment_view_test.go b/services/agents-api/internal/store/sandbox_deployment_view_test.go index b350bda1b..ec9ef0054 100644 --- a/services/agents-api/internal/store/sandbox_deployment_view_test.go +++ b/services/agents-api/internal/store/sandbox_deployment_view_test.go @@ -40,7 +40,7 @@ func TestSandboxDeploymentViewRecordsTemplateBuildAndSuspension(t *testing.T) { for _, want := range []string{ `"specification":{"resources":{"cpus":2,"memory_mib":2048}}`, `"template_build":{"status":"ready","resources":{"cpus":2,"memory_mib":2048,"root_disk_mib":24063}}`, - `"suspension":null`, + `"suspension":{"idle_seconds":300,"retention_seconds":86400}`, } { if !bytes.Contains(raw, []byte(want)) { t.Fatalf("E2B view lacks %s: %s", want, raw) @@ -73,8 +73,21 @@ func TestSandboxDeploymentViewRecordsTemplateBuildAndSuspension(t *testing.T) { if err != nil || view.Generation != 1 || !bytes.Contains(raw, recorded) { t.Fatalf("identical PUT did not record the build: %s %v", raw, err) } + // A deployment saved by an older release keeps its zero policy until the + // administrator explicitly resubmits the same E2B selection. + if _, err := pool.Exec(t.Context(), "UPDATE runtime_deployment SET idle_seconds=0,retention_seconds=0 WHERE provider_kind='e2b'"); err != nil { + t.Fatal(err) + } + view, err = s.GetRuntimeDeployment(t.Context()) + if err != nil || view.Suspension != nil { + t.Fatal("older deployment policy changed without a write", err) + } + view, err = w.UpdateSandboxDeployment(SandboxResetTestContext(t.Context()), id, SandboxDeploymentUpdateRequest{SandboxDeploymentSetupRequest: input, ExpectedGeneration: 1}) + if err != nil || view.Generation != 2 || view.Suspension == nil || view.Suspension.IdleSeconds != 300 { + t.Fatal("explicit PUT did not enable E2B idle pause", view, err) + } update := SandboxDeploymentUpdateRequest{SandboxDeploymentSetupRequest: SandboxDeploymentSetupRequest{ - DeploymentSpec: SandboxDeploymentTestSpec("microsandbox"), Provider: "microsandbox"}, ExpectedGeneration: 1} + DeploymentSpec: SandboxDeploymentTestSpec("microsandbox"), Provider: "microsandbox"}, ExpectedGeneration: 2} view, err = resetAndSelect(t, w, id, update.ExpectedGeneration, update.SandboxDeploymentSetupRequest) if err != nil || view.E2B != nil || view.Suspension == nil || view.Suspension.IdleSeconds != 300 || view.Suspension.RetentionSeconds != 86400 { t.Fatalf("microsandbox suspension view = %+v %v", view, err) diff --git a/services/agents-api/tools/e2b-provider/README.md b/services/agents-api/tools/e2b-provider/README.md index 011f49dda..8707aeeea 100644 --- a/services/agents-api/tools/e2b-provider/README.md +++ b/services/agents-api/tools/e2b-provider/README.md @@ -132,7 +132,13 @@ cleanup. Initializer and three-harness qualification remain deployment checks. Core-managed E2B is a separate hosted deployment choice, using the official pinned Python SDK through a packaged private helper. Do not restore the retired custom HTTP/Connect or envd implementation. The adapter implements the same five operations; -cloud allocations use direct placement with no synthetic node, while Runtime execution +it also implements resident memory pause/resume for Core's five-minute idle policy +when the deployment has that policy enabled. +The helper persists an outstanding pause before calling the SDK, observes an +unknown result without replaying Pause, and reconnects only the same sandbox ID +with `on_resume="restore"`. The Core allocation retains the sandbox for up to +24 hours while paused; provider snapshot storage may still be billed. +Cloud allocations use direct placement with no synthetic node, while Runtime execution and file access keep the shared daemon contract. Its immutable Runtime template build is deployment configuration, not a public Environment Template. Keep the account API key encrypted in the database, write-only through admin input and absent from helper diff --git a/services/agents-api/tools/e2b-provider/provider.py b/services/agents-api/tools/e2b-provider/provider.py index d50c06d20..e150df9e4 100644 --- a/services/agents-api/tools/e2b-provider/provider.py +++ b/services/agents-api/tools/e2b-provider/provider.py @@ -68,7 +68,7 @@ def __init__(self, request): self.reference = request['Reference'] self.references = request.get('References') or [] if (request['Version'] != 1 or request['Operation'] not in - ('create', 'inspect', 'renew', 'kill', 'command', 'validate_deployment', 'observe', 'list_templates', 'list_builds', 'verify_credential') or + ('create', 'inspect', 'renew', 'pause', 'resume', 'kill', 'command', 'validate_deployment', 'observe', 'list_templates', 'list_builds', 'verify_credential') or (request['Operation'] not in ('validate_deployment', 'observe', 'list_templates', 'list_builds', 'verify_credential') and not valid_reference(self.reference)) or (request['Operation'] == 'observe' and @@ -247,6 +247,51 @@ def renew(self): Sandbox.set_timeout(cloud.sandbox_id, self.config['TimeoutSeconds'], **self.options()) return self.qualified(self.owns(Sandbox.get_info(cloud.sandbox_id, **self.options()))) + def pause(self): + cloud = self.inspect() + if cloud is None or not self.receipt.data.get('bootstrap_complete'): + raise Failure('unconfirmed') + # A pending native call may still complete after the caller loses its + # response. Never issue a second pause that could race a later resume. + if self.receipt.data.get('pause_status') == 'pending': + if cloud.state != 'paused': + raise Failure('unconfirmed') + elif cloud.state == 'paused': + if self.receipt.data.get('pause_status') != 'settled': + raise Failure('unconfirmed') + elif cloud.state == 'running': + self.receipt.save(pause_status='pending') + Sandbox.pause(cloud.sandbox_id, keep_memory=True, **self.options()) + cloud = self.inspect() + else: + raise Failure('unconfirmed') + if cloud.state != 'paused': + raise Failure('unconfirmed') + self.receipt.save(pause_status='settled') + return cloud + + def resume(self): + cloud = self.inspect() + status = (self.receipt.data or {}).get('pause_status') + if cloud is None or cloud.state not in ('paused', 'running') or not self.receipt.data.get('bootstrap_complete') or status not in ('settled', 'resumed'): + raise Failure('unconfirmed') + if status == 'resumed': + if cloud.state != 'running': + raise Failure('unconfirmed') + return cloud + # Connect is idempotent for a running instance. It also returns fresh + # envd connection material after a paused instance is restored. + connected = Sandbox.connect(cloud.sandbox_id, timeout=self.config['TimeoutSeconds'], + on_resume='restore', **self.options()) + if connected.sandbox_id != cloud.sandbox_id: + raise Failure('ownership') + self.check_domain(connected) + detail = self.qualified(self.owns(Sandbox.get_info(cloud.sandbox_id, **self.options()))) + if detail.state != 'running': + raise Failure('unconfirmed') + self.receipt.save(connection=connection_material(connected), pause_status='resumed') + return detail + def kill(self): if self.rejected_absence(): return @@ -405,7 +450,8 @@ def execute(self): raise Failure('unconfirmed') result = run(self.client(cloud), self.q['Command'], self.remaining) return {'Version': 1, 'Command': result, 'ErrorCode': ''} - cloud = {'create': self.create, 'inspect': self.inspect, 'renew': self.renew}[operation]() + cloud = {'create': self.create, 'inspect': self.inspect, 'renew': self.renew, + 'pause': self.pause, 'resume': self.resume}[operation]() return {'Version': 1, 'Info': self.info(cloud, absent=cloud is None), 'ErrorCode': ''} except Failure as error: info = self.info() diff --git a/services/agents-api/tools/e2b-provider/provider_test.py b/services/agents-api/tools/e2b-provider/provider_test.py index 74d708547..e7f7b71a9 100644 --- a/services/agents-api/tools/e2b-provider/provider_test.py +++ b/services/agents-api/tools/e2b-provider/provider_test.py @@ -83,6 +83,38 @@ def test_create_recover_and_never_replay(self): self.assertEqual(self.record()['connection']['envd_access_token'], 'private-envd-secret') self.api.connect.assert_not_called() + def test_memory_pause_and_same_instance_resume(self): + self.assertEqual(self.call('create')['ErrorCode'], '') + def pause(sandbox_id, **kwargs): + self.assertEqual(sandbox_id, self.cloud.sandbox_id) + self.assertTrue(kwargs['keep_memory']) + self.cloud.state = 'paused' + return True + def connect(sandbox_id, **kwargs): + self.assertEqual(sandbox_id, self.cloud.sandbox_id) + self.assertEqual(kwargs['on_resume'], 'restore') + self.cloud.state = 'running' + return self.cloud + self.api.pause.side_effect = pause + self.api.connect.side_effect = connect + self.assertEqual(self.call('pause')['Info']['State'], 'paused') + self.assertEqual(self.call('pause')['Info']['State'], 'paused') + self.api.pause.assert_called_once() + self.assertEqual(self.call('resume')['Info']['State'], 'running') + self.assertEqual(self.call('resume')['Info']['State'], 'running') + self.api.connect.assert_called_once() + self.assertEqual(self.record()['ids'], [self.cloud.sandbox_id]) + + def test_unknown_pause_is_observed_without_replay(self): + self.assertEqual(self.call('create')['ErrorCode'], '') + self.api.pause.side_effect = TimeoutError('unknown pause result') + self.assertEqual(self.call('pause')['ErrorCode'], 'unconfirmed') + self.assertEqual(self.call('pause')['ErrorCode'], 'unconfirmed') + self.api.pause.assert_called_once() + self.cloud.state = 'paused' + self.assertEqual(self.call('pause')['Info']['State'], 'paused') + self.assertEqual(self.record()['pause_status'], 'settled') + def test_create_response_without_metadata_reads_exact_id_before_bootstrap(self): created = SimpleNamespace(sandbox_id=self.cloud.sandbox_id, sandbox_domain=self.cloud.sandbox_domain, From 78658bcbddfb571a2f96578c7918ac028221a180 Mon Sep 17 00:00:00 2001 From: sam Date: Wed, 30 Sep 2026 10:05:37 +0800 Subject: [PATCH 02/14] Clarify lifecycle fixture failure stage --- .../internal/store/runtime_compute_lifecycle_test.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/services/agents-api/internal/store/runtime_compute_lifecycle_test.go b/services/agents-api/internal/store/runtime_compute_lifecycle_test.go index bdc98a255..d24c760ee 100644 --- a/services/agents-api/internal/store/runtime_compute_lifecycle_test.go +++ b/services/agents-api/internal/store/runtime_compute_lifecycle_test.go @@ -325,13 +325,13 @@ func (f *computeLifecycleFixture) create() (string, store.Session, store.Environ tenant, session, environment := managedSession(t, f.store) owner, err := f.worker.ProvisionEnvironment(t.Context(), tenant, environment.ID, f.key) if err != nil { - t.Fatal(err) + t.Fatalf("provision environment: %v", err) } f.provider.mu.Lock() b := f.provider.bootstraps[owner.ID] f.provider.mu.Unlock() if err := f.provider.connect(t.Context(), b); err != nil { - t.Fatal(err) + t.Fatalf("connect runtime: %v", err) } owner = f.phase(tenant, environment.ID, "running") return tenant, session, environment, owner From 78a2c34c4ee4cd2a14f2ca7e18dcb6b7fe5e28c1 Mon Sep 17 00:00:00 2001 From: sam Date: Wed, 30 Sep 2026 10:15:23 +0800 Subject: [PATCH 03/14] Identify resident fixture failure stage --- .../internal/store/runtime_compute_lifecycle_test.go | 6 ++++-- .../agents-api/internal/store/runtime_lifecycle_test.go | 4 ++-- 2 files changed, 6 insertions(+), 4 deletions(-) diff --git a/services/agents-api/internal/store/runtime_compute_lifecycle_test.go b/services/agents-api/internal/store/runtime_compute_lifecycle_test.go index d24c760ee..d19c95e28 100644 --- a/services/agents-api/internal/store/runtime_compute_lifecycle_test.go +++ b/services/agents-api/internal/store/runtime_compute_lifecycle_test.go @@ -323,6 +323,7 @@ func (f *computeLifecycleFixture) create() (string, store.Session, store.Environ t := f.t t.Helper() tenant, session, environment := managedSession(t, f.store) + t.Log("managed session created") owner, err := f.worker.ProvisionEnvironment(t.Context(), tenant, environment.ID, f.key) if err != nil { t.Fatalf("provision environment: %v", err) @@ -334,6 +335,7 @@ func (f *computeLifecycleFixture) create() (string, store.Session, store.Environ t.Fatalf("connect runtime: %v", err) } owner = f.phase(tenant, environment.ID, "running") + t.Log("managed runtime running") return tenant, session, environment, owner } func (f *computeLifecycleFixture) phase(tenant, environment, phase string) store.RuntimeAllocation { @@ -344,11 +346,11 @@ func (f *computeLifecycleFixture) phase(tenant, environment, phase string) store err := f.worker.ReconcileManagedRuntimes(ctx) cancel() if err != nil { - f.t.Fatal(err) + f.t.Fatalf("reconcile compute phase %s: %v", phase, err) } owner, err = f.store.GetRuntimeAllocation(f.t.Context(), tenant, environment) if err != nil { - f.t.Fatal(err) + f.t.Fatalf("read compute phase %s: %v", phase, err) } if owner.ComputePhase == phase { return owner diff --git a/services/agents-api/internal/store/runtime_lifecycle_test.go b/services/agents-api/internal/store/runtime_lifecycle_test.go index b8374d62a..f5612cc80 100644 --- a/services/agents-api/internal/store/runtime_lifecycle_test.go +++ b/services/agents-api/internal/store/runtime_lifecycle_test.go @@ -120,11 +120,11 @@ func managedSession(t *testing.T, s *store.Store) (string, store.Session, store. tenant := uuid.NewString() v, e := s.CreateSession(t.Context(), tenant, store.WithFixtureModelProvider(store.CreateSessionInput{Creator: store.FixtureCreator(), Engine: "codex", IdempotencyKey: uuid.NewString(), Configuration: json.RawMessage(`{"agent":{"model":"test"},"environment":{"type":"openai_hosted","network":{"access":"enabled"}}}`)})) if e != nil { - t.Fatal(e) + t.Fatalf("create managed session: %v", e) } env, e := s.GetSessionEnvironment(t.Context(), tenant, v.ID) if e != nil { - t.Fatal(e) + t.Fatalf("load managed session environment: %v", e) } return tenant, v, env } From 416aa2cf69ed89e896db1d67d2ba24c7915d1796 Mon Sep 17 00:00:00 2001 From: sam Date: Wed, 30 Sep 2026 10:26:43 +0800 Subject: [PATCH 04/14] Initialize resident lifecycle fixture through Web deployment --- .../store/runtime_compute_lifecycle_test.go | 23 ++++++++++++++++++- 1 file changed, 22 insertions(+), 1 deletion(-) diff --git a/services/agents-api/internal/store/runtime_compute_lifecycle_test.go b/services/agents-api/internal/store/runtime_compute_lifecycle_test.go index d19c95e28..fc3dd6d19 100644 --- a/services/agents-api/internal/store/runtime_compute_lifecycle_test.go +++ b/services/agents-api/internal/store/runtime_compute_lifecycle_test.go @@ -300,7 +300,19 @@ func (f *computeLifecycleFixture) start() { t.Helper() config := &execution.RuntimeProvider{CoreURL: "http://core.invalid/api/v1", InstallationID: f.key, BackendFingerprint: "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", Provider: f.managed, Suspension: &f.policy} if _, resident := f.managed.(sandbox.ResidentPauseProvider); resident { - config.Mode, config.ProviderKind = "direct", "e2b" + f.store.SetPublicURL("https://core.example") + build := func(setup store.SandboxSetup) *execution.RuntimeProvider { + return &execution.RuntimeProvider{CoreURL: "https://core.example/api/v1", InstallationID: setup.InstallationID, BackendFingerprint: setup.BackendFingerprint, Provider: f.managed, Mode: setup.Mode, ProviderKind: setup.Provider, Generation: setup.Generation, Suspension: &f.policy} + } + config = execution.NewDeferredRuntimeProvider(f.key, func(ctx context.Context) (*execution.RuntimeProvider, error) { + setup, err := f.store.GetSandboxSetup(ctx) + if err != nil || setup.Provider == "" { + return nil, err + } + return build(setup), nil + }, func(_ context.Context, setup store.SandboxSetup) (execution.PreparedRuntimeDeployment, error) { + return execution.PreparedRuntimeDeployment{Config: build(setup)}, nil + }) } w, err := execution.StartWorker(t.Context(), &execution.Dispatcher{Store: f.store, Registry: f.provider.registry, ManagedRuntimes: config}) if err != nil { @@ -312,6 +324,15 @@ func (f *computeLifecycleFixture) start() { } f.worker, f.stop = w, stop t.Cleanup(stop) + if _, resident := f.managed.(sandbox.ResidentPauseProvider); resident { + _, err := w.InitializeSandboxDeployment(t.Context(), store.SandboxDeploymentSetupRequest{ + Provider: "e2b", DeploymentSpec: store.SandboxDeploymentTestSpec("e2b"), + E2B: &sandbox.E2BConfiguration{APIKey: "fixture-api-key", Template: "runtime:" + uuid.NewString(), TemplateBuild: &sandbox.TemplateBuild{Status: "ready", CPUs: 2, MemoryMiB: 2048}}, + }) + if err != nil { + t.Fatalf("initialize resident test deployment: %v", err) + } + } } func (f *computeLifecycleFixture) sql(query string, args ...any) { f.t.Helper() From 0b95676cb0f9b76d01f15b893424dd1b4997a6e8 Mon Sep 17 00:00:00 2001 From: sam Date: Wed, 30 Sep 2026 10:36:58 +0800 Subject: [PATCH 05/14] Observe uncertain resident pause in lifecycle test --- .../store/runtime_compute_lifecycle_test.go | 2 -- .../store/runtime_resident_lifecycle_test.go | 13 +++++++++++-- 2 files changed, 11 insertions(+), 4 deletions(-) diff --git a/services/agents-api/internal/store/runtime_compute_lifecycle_test.go b/services/agents-api/internal/store/runtime_compute_lifecycle_test.go index fc3dd6d19..c015849b4 100644 --- a/services/agents-api/internal/store/runtime_compute_lifecycle_test.go +++ b/services/agents-api/internal/store/runtime_compute_lifecycle_test.go @@ -344,7 +344,6 @@ func (f *computeLifecycleFixture) create() (string, store.Session, store.Environ t := f.t t.Helper() tenant, session, environment := managedSession(t, f.store) - t.Log("managed session created") owner, err := f.worker.ProvisionEnvironment(t.Context(), tenant, environment.ID, f.key) if err != nil { t.Fatalf("provision environment: %v", err) @@ -356,7 +355,6 @@ func (f *computeLifecycleFixture) create() (string, store.Session, store.Environ t.Fatalf("connect runtime: %v", err) } owner = f.phase(tenant, environment.ID, "running") - t.Log("managed runtime running") return tenant, session, environment, owner } func (f *computeLifecycleFixture) phase(tenant, environment, phase string) store.RuntimeAllocation { diff --git a/services/agents-api/internal/store/runtime_resident_lifecycle_test.go b/services/agents-api/internal/store/runtime_resident_lifecycle_test.go index e76c77725..fc60915ee 100644 --- a/services/agents-api/internal/store/runtime_resident_lifecycle_test.go +++ b/services/agents-api/internal/store/runtime_resident_lifecycle_test.go @@ -107,8 +107,17 @@ func TestResidentUnknownPauseUsesObservationWithoutReplay(t *testing.T) { tenant, _, environment, owner := f.create() f.phase(tenant, environment.ID, "running") f.complete(owner) - if err := f.worker.ReconcileManagedRuntimes(t.Context()); !errors.Is(err, sandbox.ErrComputeUnconfirmed) { - t.Fatal("unknown pause was reported as settled", err) + for range 100 { + if err := f.worker.ReconcileManagedRuntimes(t.Context()); err != nil && !errors.Is(err, sandbox.ErrComputeUnconfirmed) { + t.Fatal("resident pause reconciliation failed", err) + } + if resident.pauses == 1 { + break + } + } + uncertain, err := f.store.GetRuntimeAllocation(t.Context(), tenant, environment.ID) + if err != nil || resident.pauses != 1 || uncertain.ComputePhase != "suspending" { + t.Fatal("unknown pause was not retained for observation", resident.pauses, uncertain.ComputePhase, err) } f.phase(tenant, environment.ID, "suspended") if resident.pauses != 1 { From 7062dfb3ac36a8c155f57a59ecb17ed805c8daec Mon Sep 17 00:00:00 2001 From: sam Date: Wed, 30 Sep 2026 10:48:17 +0800 Subject: [PATCH 06/14] Observe uncertain resident pause and resume without replay --- .../execution/runtime_compute_resident.go | 21 +++++++++--- .../store/runtime_resident_lifecycle_test.go | 34 +++++++++++++++++++ 2 files changed, 50 insertions(+), 5 deletions(-) diff --git a/services/agents-api/internal/execution/runtime_compute_resident.go b/services/agents-api/internal/execution/runtime_compute_resident.go index a987f1cc6..502c8f1b3 100644 --- a/services/agents-api/internal/execution/runtime_compute_resident.go +++ b/services/agents-api/internal/execution/runtime_compute_resident.go @@ -44,7 +44,10 @@ func (r *runtimeLifecycle) observeResidentCompute(ctx context.Context, p sandbox } return r.wakeResidentCompute(ctx, p, next, state) case "suspending": - info, err := p.Pause(ctx, runtimeReference(owner)) + // The pause receipt precedes the external call. After an uncertain + // result, even a running observation cannot prove that the original + // request will not pause later, so never issue Pause a second time. + info, err := p.GetInfo(ctx, runtimeReference(owner)) if err != nil { return err } @@ -66,9 +69,9 @@ func (r *runtimeLifecycle) observeResidentCompute(ctx context.Context, p sandbox if err != nil { return err } - return r.restoreResidentCompute(ctx, p, next, state) + return r.restoreResidentCompute(ctx, p, next, state, true) case "restoring": - return r.restoreResidentCompute(ctx, p, owner, state) + return r.restoreResidentCompute(ctx, p, owner, state, false) case "waking": return r.wakeResidentCompute(ctx, p, owner, state) default: @@ -160,8 +163,16 @@ func (r *runtimeLifecycle) idleResidentCompute(ctx context.Context, p sandbox.Re return err } -func (r *runtimeLifecycle) restoreResidentCompute(ctx context.Context, p sandbox.ResidentPauseProvider, owner store.RuntimeAllocation, state runtimeCompute) error { - info, err := p.Resume(ctx, runtimeReference(owner)) +func (r *runtimeLifecycle) restoreResidentCompute(ctx context.Context, p sandbox.ResidentPauseProvider, owner store.RuntimeAllocation, state runtimeCompute, issue bool) error { + var info sandbox.Info + var err error + if issue { + info, err = p.Resume(ctx, runtimeReference(owner)) + } else { + // A lost Resume response may have restored the original VM. Observe + // the durable restoring receipt without sending another Resume. + info, err = p.GetInfo(ctx, runtimeReference(owner)) + } if err != nil { return err } diff --git a/services/agents-api/internal/store/runtime_resident_lifecycle_test.go b/services/agents-api/internal/store/runtime_resident_lifecycle_test.go index fc60915ee..1e28213e4 100644 --- a/services/agents-api/internal/store/runtime_resident_lifecycle_test.go +++ b/services/agents-api/internal/store/runtime_resident_lifecycle_test.go @@ -13,6 +13,7 @@ type fakeResidentProvider struct { *fakeCheckpointProvider pauses, resumes int losePause bool + loseResume bool } func (p *fakeResidentProvider) Create(ctx context.Context, b sandbox.Bootstrap) (sandbox.Info, error) { @@ -58,6 +59,10 @@ func (p *fakeResidentProvider) Resume(ctx context.Context, r sandbox.Reference) p.resumes++ info.State = "running" p.resources[r.AllocationID] = info + if p.loseResume { + p.loseResume = false + return sandbox.Info{}, sandbox.ErrComputeUnconfirmed + } return info, nil } func (p *fakeResidentProvider) RunCommand(ctx context.Context, r sandbox.Reference, command sandbox.Command) (sandbox.CommandResult, error) { @@ -125,6 +130,35 @@ func TestResidentUnknownPauseUsesObservationWithoutReplay(t *testing.T) { } } +func TestResidentUnknownResumeUsesObservationWithoutReplay(t *testing.T) { + var resident *fakeResidentProvider + f := newComputeLifecycleFixtureWithProvider(t, 2, 4, func(p *fakeCheckpointProvider) sandbox.SandboxProvider { + resident = &fakeResidentProvider{fakeCheckpointProvider: p, loseResume: true} + return resident + }) + + tenant, _, environment, owner := f.create() + f.complete(owner) + owner = f.phase(tenant, environment.ID, "suspended") + f.queued(owner) + for range 100 { + if err := f.worker.ReconcileManagedRuntimes(t.Context()); err != nil && !errors.Is(err, sandbox.ErrComputeUnconfirmed) { + t.Fatal("resident resume reconciliation failed", err) + } + if resident.resumes == 1 { + break + } + } + uncertain, err := f.store.GetRuntimeAllocation(t.Context(), tenant, environment.ID) + if err != nil || resident.resumes != 1 || uncertain.ComputePhase != "restoring" { + t.Fatal("unknown resume was not retained for observation", resident.resumes, uncertain.ComputePhase, err) + } + f.phase(tenant, environment.ID, "running") + if resident.resumes != 1 { + t.Fatal("unknown resume was replayed", resident.resumes) + } +} + func TestResidentDoesNotPauseWhileTurnQueued(t *testing.T) { var resident *fakeResidentProvider f := newComputeLifecycleFixtureWithProvider(t, 2, 4, func(p *fakeCheckpointProvider) sandbox.SandboxProvider { From 363fcc87e6c40134bb3f99d8672b02ac0e3e5f67 Mon Sep 17 00:00:00 2001 From: sam Date: Wed, 30 Sep 2026 10:52:02 +0800 Subject: [PATCH 07/14] Keep idempotent E2B reconnect on uncertain resume --- .../execution/runtime_compute_resident.go | 16 +++------ .../store/runtime_resident_lifecycle_test.go | 34 ------------------- 2 files changed, 4 insertions(+), 46 deletions(-) diff --git a/services/agents-api/internal/execution/runtime_compute_resident.go b/services/agents-api/internal/execution/runtime_compute_resident.go index 502c8f1b3..fce5fdb82 100644 --- a/services/agents-api/internal/execution/runtime_compute_resident.go +++ b/services/agents-api/internal/execution/runtime_compute_resident.go @@ -69,9 +69,9 @@ func (r *runtimeLifecycle) observeResidentCompute(ctx context.Context, p sandbox if err != nil { return err } - return r.restoreResidentCompute(ctx, p, next, state, true) + return r.restoreResidentCompute(ctx, p, next, state) case "restoring": - return r.restoreResidentCompute(ctx, p, owner, state, false) + return r.restoreResidentCompute(ctx, p, owner, state) case "waking": return r.wakeResidentCompute(ctx, p, owner, state) default: @@ -163,16 +163,8 @@ func (r *runtimeLifecycle) idleResidentCompute(ctx context.Context, p sandbox.Re return err } -func (r *runtimeLifecycle) restoreResidentCompute(ctx context.Context, p sandbox.ResidentPauseProvider, owner store.RuntimeAllocation, state runtimeCompute, issue bool) error { - var info sandbox.Info - var err error - if issue { - info, err = p.Resume(ctx, runtimeReference(owner)) - } else { - // A lost Resume response may have restored the original VM. Observe - // the durable restoring receipt without sending another Resume. - info, err = p.GetInfo(ctx, runtimeReference(owner)) - } +func (r *runtimeLifecycle) restoreResidentCompute(ctx context.Context, p sandbox.ResidentPauseProvider, owner store.RuntimeAllocation, state runtimeCompute) error { + info, err := p.Resume(ctx, runtimeReference(owner)) if err != nil { return err } diff --git a/services/agents-api/internal/store/runtime_resident_lifecycle_test.go b/services/agents-api/internal/store/runtime_resident_lifecycle_test.go index 1e28213e4..fc60915ee 100644 --- a/services/agents-api/internal/store/runtime_resident_lifecycle_test.go +++ b/services/agents-api/internal/store/runtime_resident_lifecycle_test.go @@ -13,7 +13,6 @@ type fakeResidentProvider struct { *fakeCheckpointProvider pauses, resumes int losePause bool - loseResume bool } func (p *fakeResidentProvider) Create(ctx context.Context, b sandbox.Bootstrap) (sandbox.Info, error) { @@ -59,10 +58,6 @@ func (p *fakeResidentProvider) Resume(ctx context.Context, r sandbox.Reference) p.resumes++ info.State = "running" p.resources[r.AllocationID] = info - if p.loseResume { - p.loseResume = false - return sandbox.Info{}, sandbox.ErrComputeUnconfirmed - } return info, nil } func (p *fakeResidentProvider) RunCommand(ctx context.Context, r sandbox.Reference, command sandbox.Command) (sandbox.CommandResult, error) { @@ -130,35 +125,6 @@ func TestResidentUnknownPauseUsesObservationWithoutReplay(t *testing.T) { } } -func TestResidentUnknownResumeUsesObservationWithoutReplay(t *testing.T) { - var resident *fakeResidentProvider - f := newComputeLifecycleFixtureWithProvider(t, 2, 4, func(p *fakeCheckpointProvider) sandbox.SandboxProvider { - resident = &fakeResidentProvider{fakeCheckpointProvider: p, loseResume: true} - return resident - }) - - tenant, _, environment, owner := f.create() - f.complete(owner) - owner = f.phase(tenant, environment.ID, "suspended") - f.queued(owner) - for range 100 { - if err := f.worker.ReconcileManagedRuntimes(t.Context()); err != nil && !errors.Is(err, sandbox.ErrComputeUnconfirmed) { - t.Fatal("resident resume reconciliation failed", err) - } - if resident.resumes == 1 { - break - } - } - uncertain, err := f.store.GetRuntimeAllocation(t.Context(), tenant, environment.ID) - if err != nil || resident.resumes != 1 || uncertain.ComputePhase != "restoring" { - t.Fatal("unknown resume was not retained for observation", resident.resumes, uncertain.ComputePhase, err) - } - f.phase(tenant, environment.ID, "running") - if resident.resumes != 1 { - t.Fatal("unknown resume was replayed", resident.resumes) - } -} - func TestResidentDoesNotPauseWhileTurnQueued(t *testing.T) { var resident *fakeResidentProvider f := newComputeLifecycleFixtureWithProvider(t, 2, 4, func(p *fakeCheckpointProvider) sandbox.SandboxProvider { From 6c141fabb59faa5c2936c8217019f07299e8153d Mon Sep 17 00:00:00 2001 From: sam Date: Wed, 30 Sep 2026 11:04:00 +0800 Subject: [PATCH 08/14] Keep legacy downgrade tests on zero E2B idle policy --- .../store/archived_cancellation_migration_test.go | 1 + .../internal/store/e2b_idle_downgrade_test.go | 15 +++++++++++++++ .../store/provider_registration_migration_test.go | 1 + .../internal/store/sandbox_generations_test.go | 1 + 4 files changed, 18 insertions(+) create mode 100644 services/agents-api/internal/store/e2b_idle_downgrade_test.go diff --git a/services/agents-api/internal/store/archived_cancellation_migration_test.go b/services/agents-api/internal/store/archived_cancellation_migration_test.go index 0ffa52a2d..e5ec38788 100644 --- a/services/agents-api/internal/store/archived_cancellation_migration_test.go +++ b/services/agents-api/internal/store/archived_cancellation_migration_test.go @@ -23,6 +23,7 @@ func TestArchivedCancellationMigrationDoesNotAdoptOldRevocations(t *testing.T) { } db := sql.OpenDB(stdlib.GetConnector(*s.pool.Config().ConnConfig)) defer db.Close() + disableE2BIdlePolicyForLegacyDowngrade(t, db) provider, err := goose.NewProvider(goose.DialectPostgres, db, os.DirFS("../../migrations"), goose.WithTableName("agents_api_schema_version")) if err != nil { t.Fatal(err) diff --git a/services/agents-api/internal/store/e2b_idle_downgrade_test.go b/services/agents-api/internal/store/e2b_idle_downgrade_test.go new file mode 100644 index 000000000..693d9e4cf --- /dev/null +++ b/services/agents-api/internal/store/e2b_idle_downgrade_test.go @@ -0,0 +1,15 @@ +package store + +import ( + "database/sql" + "testing" +) + +// These tests exercise older migration guards. Their historical E2B schema +// accepts only the zero idle policy, so opt the isolated fixture out first. +func disableE2BIdlePolicyForLegacyDowngrade(t *testing.T, db *sql.DB) { + t.Helper() + if _, err := db.ExecContext(t.Context(), `UPDATE runtime_deployment SET idle_seconds=0, retention_seconds=0 WHERE provider_kind='e2b'`); err != nil { + t.Fatal(err) + } +} diff --git a/services/agents-api/internal/store/provider_registration_migration_test.go b/services/agents-api/internal/store/provider_registration_migration_test.go index 1b539a98a..ffee167f9 100644 --- a/services/agents-api/internal/store/provider_registration_migration_test.go +++ b/services/agents-api/internal/store/provider_registration_migration_test.go @@ -27,6 +27,7 @@ func TestProviderRegistrationDowngradePreservesCustomEndpoints(t *testing.T) { } db := sql.OpenDB(stdlib.GetConnector(*s.pool.Config().ConnConfig)) defer db.Close() + disableE2BIdlePolicyForLegacyDowngrade(t, db) migrations, err := goose.NewProvider(goose.DialectPostgres, db, os.DirFS("../../migrations"), goose.WithTableName("agents_api_schema_version")) if err != nil { t.Fatal(err) diff --git a/services/agents-api/internal/store/sandbox_generations_test.go b/services/agents-api/internal/store/sandbox_generations_test.go index 09d749eed..90a982345 100644 --- a/services/agents-api/internal/store/sandbox_generations_test.go +++ b/services/agents-api/internal/store/sandbox_generations_test.go @@ -287,6 +287,7 @@ func TestGenerationDowngradeRefusesOldAllocation(t *testing.T) { } db := sql.OpenDB(stdlib.GetConnector(*s.pool.Config().ConnConfig)) defer db.Close() + disableE2BIdlePolicyForLegacyDowngrade(t, db) migrations, err := goose.NewProvider(goose.DialectPostgres, db, os.DirFS("../../migrations"), goose.WithTableName("agents_api_schema_version")) if err != nil { t.Fatal(err) From 52123a35c65381fee17bee726c75380eba8264ab Mon Sep 17 00:00:00 2001 From: sam Date: Wed, 30 Sep 2026 13:49:58 +0800 Subject: [PATCH 09/14] Preserve resident pause capability through generation routing --- .../cmd/server/managed_generations.go | 11 ++++- .../server/managed_generations_resident.go | 48 +++++++++++++++++++ .../cmd/server/managed_generations_test.go | 32 +++++++++++-- 3 files changed, 85 insertions(+), 6 deletions(-) create mode 100644 services/agents-api/cmd/server/managed_generations_resident.go diff --git a/services/agents-api/cmd/server/managed_generations.go b/services/agents-api/cmd/server/managed_generations.go index cb3ba82c2..8b3cf3331 100644 --- a/services/agents-api/cmd/server/managed_generations.go +++ b/services/agents-api/cmd/server/managed_generations.go @@ -113,9 +113,16 @@ func (s *managedSetup) routeGenerations(candidate execution.PreparedRuntimeDeplo return execution.PreparedRuntimeDeployment{}, errors.New("sandbox generation store is unavailable") } router := &generationRouter{setup: s, store: db} - if _, ok := candidate.Config.Provider.(runtimeobs.Source); ok { + _, observed := candidate.Config.Provider.(runtimeobs.Source) + _, resident := candidate.Config.Provider.(sandbox.ResidentPauseProvider) + switch { + case observed && resident: + candidate.Config.Provider = &observedResidentGenerationRouter{&residentGenerationRouter{router}} + case resident: + candidate.Config.Provider = &residentGenerationRouter{router} + case observed: candidate.Config.Provider = &observedGenerationRouter{router} - } else { + default: candidate.Config.Provider = router } if !adapter.Credential { diff --git a/services/agents-api/cmd/server/managed_generations_resident.go b/services/agents-api/cmd/server/managed_generations_resident.go new file mode 100644 index 000000000..b8029ef07 --- /dev/null +++ b/services/agents-api/cmd/server/managed_generations_resident.go @@ -0,0 +1,48 @@ +package main + +import ( + "context" + + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeobs" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" +) + +// Only providers that advertise resident pause receive this facade. Every +// operation keeps the allocation's immutable generation and credential fence. +type residentGenerationRouter struct{ *generationRouter } + +func (p *residentGenerationRouter) Pause(ctx context.Context, ref sandbox.Reference) (sandbox.Info, error) { + provider, done, err := p.route(ctx, ref) + if err != nil { + return sandbox.Info{}, err + } + defer done() + resident, ok := provider.(sandbox.ResidentPauseProvider) + if !ok { + return sandbox.Info{}, sandbox.ErrInvalid + } + return resident.Pause(ctx, ref) +} + +func (p *residentGenerationRouter) Resume(ctx context.Context, ref sandbox.Reference) (sandbox.Info, error) { + provider, done, err := p.route(ctx, ref) + if err != nil { + return sandbox.Info{}, err + } + defer done() + resident, ok := provider.(sandbox.ResidentPauseProvider) + if !ok { + return sandbox.Info{}, sandbox.ErrInvalid + } + return resident.Resume(ctx, ref) +} + +type observedResidentGenerationRouter struct{ *residentGenerationRouter } + +func (p *observedResidentGenerationRouter) Observe(ctx context.Context, target runtimeobs.Target) (runtimeobs.Sample, error) { + return (&observedGenerationRouter{p.generationRouter}).Observe(ctx, target) +} + +var _ sandbox.ResidentPauseProvider = (*residentGenerationRouter)(nil) +var _ sandbox.ResidentPauseProvider = (*observedResidentGenerationRouter)(nil) +var _ runtimeobs.Source = (*observedResidentGenerationRouter)(nil) diff --git a/services/agents-api/cmd/server/managed_generations_test.go b/services/agents-api/cmd/server/managed_generations_test.go index 4d9190d0b..852c34250 100644 --- a/services/agents-api/cmd/server/managed_generations_test.go +++ b/services/agents-api/cmd/server/managed_generations_test.go @@ -10,6 +10,8 @@ import ( "testing" "time" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/execution" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeobs" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox/e2b" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" @@ -42,8 +44,9 @@ func TestE2BRouterKeepsOldSpecificationWithCommittedCredential(t *testing.T) { import json,sys,pathlib q=json.load(sys.stdin) with (pathlib.Path(q['Config']['StateDir'])/'requests').open('a') as f: f.write(json.dumps(q)+'\n') -info=dict(q['Reference'],State='running',ProviderID='owned',CreateSettled=True) +info=dict(q['Reference'],State='running',ProviderID='owned',CreateSettled=True,BootstrapComplete=True) if q['Operation']=='kill': info['State']='absent' +if q['Operation']=='pause': info['State']='paused' print(json.dumps({'Version':1,'Info':info})) ` if err := os.WriteFile(helper, []byte(script), 0700); err != nil { @@ -61,7 +64,22 @@ print(json.dumps({'Version':1,'Info':info})) db := &routingSetupStore{setupStore: setupStore{value: current}, old: old, oldID: ref.AllocationID} setup := &managedSetup{store: db, installationID: id} // A facade retained by a generation-one lifecycle still reads current credentials. - router := &generationRouter{setup: setup, store: db} + provider, err := setup.provider(current) + if err != nil { + t.Fatal(err) + } + candidate, err := setup.routeGenerations(execution.PreparedRuntimeDeployment{Config: &execution.RuntimeProvider{Provider: provider}}, current) + if err != nil { + t.Fatal(err) + } + router := candidate.Config.Provider + resident, ok := router.(sandbox.ResidentPauseProvider) + if !ok { + t.Fatal("generation facade lost resident pause capability") + } + if _, ok := router.(runtimeobs.Source); !ok { + t.Fatal("generation facade lost observation capability") + } ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second) defer cancel() if _, err := router.GetInfo(ctx, ref); err != nil { @@ -70,6 +88,12 @@ print(json.dumps({'Version':1,'Info':info})) if _, err := router.Renew(ctx, ref); err != nil { t.Fatal(err) } + if _, err := resident.Pause(ctx, ref); err != nil { + t.Fatal(err) + } + if _, err := resident.Resume(ctx, ref); err != nil { + t.Fatal(err) + } if err := router.Kill(ctx, ref); err != nil { t.Fatal(err) } @@ -83,7 +107,7 @@ print(json.dumps({'Version':1,'Info':info})) t.Fatal(err) } lines := strings.Split(strings.TrimSpace(string(raw)), "\n") - if len(lines) != 4 { + if len(lines) != 6 { t.Fatal(len(lines)) } for i, line := range lines { @@ -97,7 +121,7 @@ print(json.dumps({'Version':1,'Info':info})) t.Fatal(err) } expected := old - if i == 3 { + if i == 5 { expected = current } if q.Config.APIKey != "new-key" || q.Config.Template != expected.E2B.Template || q.Config.Resources != expected.Specification.Resources { From ca3675bf77fa1bf22698f47f890d1af2ec1f8d3b Mon Sep 17 00:00:00 2001 From: sam Date: Wed, 30 Sep 2026 14:22:30 +0800 Subject: [PATCH 10/14] fix(core): attach audit context to sandbox deployment changes --- .../internal/api/sandbox_deployment_changes_test.go | 7 ++++++- .../agents-api/internal/api/sandbox_deployment_setup.go | 1 + 2 files changed, 7 insertions(+), 1 deletion(-) diff --git a/services/agents-api/internal/api/sandbox_deployment_changes_test.go b/services/agents-api/internal/api/sandbox_deployment_changes_test.go index 5227157a1..9c5e3bd2f 100644 --- a/services/agents-api/internal/api/sandbox_deployment_changes_test.go +++ b/services/agents-api/internal/api/sandbox_deployment_changes_test.go @@ -8,6 +8,7 @@ import ( "testing" "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/device" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/adminaudit" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" ) @@ -15,8 +16,12 @@ func TestSandboxDeploymentChangesAuthenticateAndDecode(t *testing.T) { project, _ := NewAuthenticator([]APIKey{callerBinding()}) admin, _ := NewDeploymentAuthenticator([]string{device.HashCredential("administrator")}) updates, resets := 0, 0 - update := func(_ context.Context, in store.SandboxDeploymentUpdateRequest) (store.RuntimeDeploymentView, error) { + update := func(ctx context.Context, in store.SandboxDeploymentUpdateRequest) (store.RuntimeDeploymentView, error) { updates++ + source, ok := adminaudit.FromContext(ctx) + if !ok || source.ProjectID != "" || source.CredentialID == "" || source.RequestID == "" || source.TraceID == "" { + t.Fatal("deployment mutation lost administrator audit source") + } if in.Provider != "e2b" || in.ExpectedGeneration != 2 || in.E2B == nil || in.E2B.APIKey != "synthetic-private-key" { t.Fatal("write-only fields were lost") } diff --git a/services/agents-api/internal/api/sandbox_deployment_setup.go b/services/agents-api/internal/api/sandbox_deployment_setup.go index 5c06a104b..a7e5be13c 100644 --- a/services/agents-api/internal/api/sandbox_deployment_setup.go +++ b/services/agents-api/internal/api/sandbox_deployment_setup.go @@ -134,6 +134,7 @@ func (h *Handler) updateSandboxDeployment(w http.ResponseWriter, r *http.Request writeStoreError(w, r, store.ErrSandboxDeploymentConflict) return } + setAdminAuditSource(r, "") result, err := h.sandboxUpdate(r.Context(), store.SandboxDeploymentUpdateRequest{SandboxDeploymentSetupRequest: input.request(), ExpectedGeneration: *input.ExpectedGeneration}) if err != nil { writeStoreError(w, r, err) From 1e84fd293cdbb53075459c3ac28f2fe5daef56c7 Mon Sep 17 00:00:00 2001 From: sam Date: Wed, 30 Sep 2026 15:13:54 +0800 Subject: [PATCH 11/14] fix(store): compare registered sandbox suspension policies --- apps/docs/content/docs/sandbox-provider.mdx | 9 +++- apps/docs/content/guide-sources.json | 4 +- docs/sandbox-provider.md | 9 +++- .../store/sandbox_deployment_mutations.go | 18 +++---- .../store/sandbox_policy_comparison_test.go | 50 +++++++++++++++++++ 5 files changed, 76 insertions(+), 14 deletions(-) create mode 100644 services/agents-api/internal/store/sandbox_policy_comparison_test.go diff --git a/apps/docs/content/docs/sandbox-provider.mdx b/apps/docs/content/docs/sandbox-provider.mdx index 9e9f336a3..245a92209 100644 --- a/apps/docs/content/docs/sandbox-provider.mdx +++ b/apps/docs/content/docs/sandbox-provider.mdx @@ -23,7 +23,14 @@ workspaces through the common [daemon protocol](/runtime-protocol); Harness adapters translate execution. A user-owned machine uses the same Runtime contract but has no Core-owned allocation to create or destroy. -Use maintained provider SDKs behind thin adapters. Hosted deployments select one +Use maintained provider SDKs behind thin adapters. Thin means replacing a Provider +requires no changes to the common execution flow, not that an adapter contains +little code. Core owns shared scheduling, persistence and recovery through +capability contracts such as `ResidentPauseProvider`. Provider SDK calls and native +behavior belong inside the adapter; shared policy decisions use the registered +provider policy rather than vendor-name branches. + +Hosted deployments select one deployment-wide Provider: E2B cloud, or Docker/microsandbox on administrator-owned nodes. Providers never execute Core initialization commands; initialization, daily execution and Files use the daemon's typed Runtime operations. diff --git a/apps/docs/content/guide-sources.json b/apps/docs/content/guide-sources.json index 21a79edc2..3727c3986 100644 --- a/apps/docs/content/guide-sources.json +++ b/apps/docs/content/guide-sources.json @@ -28,7 +28,7 @@ "contracts/agents-api/harness-onboarding.md": "874a80dc1ffc5e97b2783bea3f6b217893f997b38bb8180a29f09c6dbf75e622", "docs/runtime-bootstrap.md": "0d49aed73b298039e04453e6f465b0e925fb2227fa35d4206820bb3acbd8df39", "docs/runtime-protocol.md": "e8aa4cf862b5f63a4138cfeb25196a6bdb5c40a2e389e828b464585415dc6c45", - "docs/sandbox-provider.md": "e19f4947bb51edd9a3a9fefb2b0a9af91adcca822f1e9b6d1066c2426c575649", + "docs/sandbox-provider.md": "a20da12b5963461128e7a7581320a3eeddd6c0eb0dbdc10e35c3ab0a9e7c12fe", "apps/docs/scripts/guides.json": "3c3768fb94fd3d42464d4ce8c55724b9ba8360133c7cc19d513d8ec83a8e034f" }, "outputs": { @@ -59,6 +59,6 @@ "content/docs/harness-onboarding.mdx": "17626e9f6256ddfffc95fd81dfa8d9c712d9ff16792ea5ef54bfe8c3ce1a2a9f", "content/docs/runtime-bootstrap.mdx": "58982955811c9ad46a9fc04d8d2aef5a762fc9d25d5833f88aa75ee3fc46523f", "content/docs/runtime-protocol.mdx": "d6b818edf2a0e0c4f2d3878f75805043c6dba3b09c6563f80faa4cecd05d40c7", - "content/docs/sandbox-provider.mdx": "71a68a37dbe443655cd7df472ff269b2bb0a065767b2e9dd32442fd57e894e2b" + "content/docs/sandbox-provider.mdx": "eeeab2d7963423928ca00e46c31e6a4d1ec56bb9ede442e3b6249f5114590fc2" } } diff --git a/docs/sandbox-provider.md b/docs/sandbox-provider.md index 0519010f6..9d89060ee 100644 --- a/docs/sandbox-provider.md +++ b/docs/sandbox-provider.md @@ -20,7 +20,14 @@ workspaces through the common [daemon protocol](runtime-protocol.md); Harness adapters translate execution. A user-owned machine uses the same Runtime contract but has no Core-owned allocation to create or destroy. -Use maintained provider SDKs behind thin adapters. Hosted deployments select one +Use maintained provider SDKs behind thin adapters. Thin means replacing a Provider +requires no changes to the common execution flow, not that an adapter contains +little code. Core owns shared scheduling, persistence and recovery through +capability contracts such as `ResidentPauseProvider`. Provider SDK calls and native +behavior belong inside the adapter; shared policy decisions use the registered +provider policy rather than vendor-name branches. + +Hosted deployments select one deployment-wide Provider: E2B cloud, or Docker/microsandbox on administrator-owned nodes. Providers never execute Core initialization commands; initialization, daily execution and Files use the daemon's typed Runtime operations. diff --git a/services/agents-api/internal/store/sandbox_deployment_mutations.go b/services/agents-api/internal/store/sandbox_deployment_mutations.go index 8ff1cc576..84e319cb5 100644 --- a/services/agents-api/internal/store/sandbox_deployment_mutations.go +++ b/services/agents-api/internal/store/sandbox_deployment_mutations.go @@ -39,16 +39,14 @@ func (s *Store) sandboxSelectionEqual(d sqlc.RuntimeDeployment, input SandboxDep if d.ProviderKind != input.Provider { return false, nil } - if input.Provider == "e2b" { - policy, err := providers.Describe(input.Provider, runtimeUUID(d.InstallationID)) - if err != nil { - return false, err - } - if d.IdleSeconds != policy.IdleSeconds || d.RetentionSeconds != policy.RetentionSeconds { - // A same-selection admin PUT explicitly opts an older deployment in - // without changing every installed E2B deployment during migration. - return false, nil - } + policy, err := providers.Describe(input.Provider, runtimeUUID(d.InstallationID)) + if err != nil { + return false, err + } + if d.IdleSeconds != policy.IdleSeconds || d.RetentionSeconds != policy.RetentionSeconds { + // An explicit configuration write adopts the registered policy without + // changing installed deployments during migration or read-only access. + return false, nil } if input.E2B == nil { return d.E2bTemplate == "", nil diff --git a/services/agents-api/internal/store/sandbox_policy_comparison_test.go b/services/agents-api/internal/store/sandbox_policy_comparison_test.go new file mode 100644 index 000000000..2f9b916d4 --- /dev/null +++ b/services/agents-api/internal/store/sandbox_policy_comparison_test.go @@ -0,0 +1,50 @@ +package store + +import ( + "encoding/json" + "testing" + + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/db/sqlc" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox/providers" + "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgtype" +) + +func TestSandboxSelectionComparesRegisteredPolicyForNodeProviders(t *testing.T) { + for _, kind := range []string{"docker", "microsandbox"} { + t.Run(kind, func(t *testing.T) { + id := uuid.New() + policy, err := providers.Describe(kind, id.String()) + if err != nil { + t.Fatal(err) + } + input := SandboxDeploymentSetupRequest{Provider: kind, DeploymentSpec: SandboxDeploymentTestSpec(kind)} + spec, err := json.Marshal(input.DeploymentSpec) + if err != nil { + t.Fatal(err) + } + original := sqlc.RuntimeDeployment{ProviderKind: kind, InstallationID: pgtype.UUID{Bytes: id, Valid: true}, Specification: spec, IdleSeconds: policy.IdleSeconds, RetentionSeconds: policy.RetentionSeconds} + for _, tc := range []struct { + name string + idle, retention int64 + equal bool + }{ + {"unchanged", policy.IdleSeconds, policy.RetentionSeconds, true}, + {"idle_policy_changed", policy.IdleSeconds + 1, policy.RetentionSeconds, false}, + {"retention_policy_changed", policy.IdleSeconds, policy.RetentionSeconds + 1, false}, + } { + t.Run(tc.name, func(t *testing.T) { + d := original + d.IdleSeconds, d.RetentionSeconds = tc.idle, tc.retention + got, err := (&Store{}).sandboxSelectionEqual(d, input) + if err != nil || got != tc.equal { + t.Fatalf("selection equality = %v, %v; want %v", got, err, tc.equal) + } + if d.IdleSeconds != tc.idle || d.RetentionSeconds != tc.retention { + t.Fatal("comparison mutated stored policy") + } + }) + } + }) + } +} From 2b9fe5ef784797b142da697ba4a82d69fe940a0b Mon Sep 17 00:00:00 2001 From: sam Date: Wed, 30 Sep 2026 15:20:16 +0800 Subject: [PATCH 12/14] fix(sandbox): refresh idle policy and enable E2B runtime suspension --- .../agents-api/deploy/e2b/managed_init.py | 7 +++ .../deploy/e2b/managed_init_test.py | 6 +++ .../internal/execution/runtime_lifecycle.go | 30 +++++++------ .../internal/execution/runtime_manager.go | 15 ++++++- .../execution/runtime_manager_test.go | 43 +++++++++++++++++++ 5 files changed, 87 insertions(+), 14 deletions(-) diff --git a/services/agents-api/deploy/e2b/managed_init.py b/services/agents-api/deploy/e2b/managed_init.py index 038a21040..615010dfb 100644 --- a/services/agents-api/deploy/e2b/managed_init.py +++ b/services/agents-api/deploy/e2b/managed_init.py @@ -13,6 +13,9 @@ _spec.loader.exec_module(shared) +SUSPEND_DIRECTORY = Path("/run/oac") + + def identity(payload): fields = ['InstallationID', 'TenantID', 'EnvironmentID', 'AllocationID', 'SessionID', 'DeviceID'] if not isinstance(payload, dict) or set(payload) != set(fields + ['NetworkAccess', 'AllowedDomains', 'RuntimeBootstrap']): @@ -42,6 +45,10 @@ def initialize(): binding = identity(payload) shared.write_private(root / 'managed-launch.json', binding) environment = shared.prepare_runtime() + SUSPEND_DIRECTORY.mkdir(mode=0o700, parents=True, exist_ok=True) + os.chown(SUSPEND_DIRECTORY, 1000, 1000) + SUSPEND_DIRECTORY.chmod(0o700) + environment['OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE'] = str(SUSPEND_DIRECTORY / 'daemon-suspend.json') environment.update(OAC_RUNTIME_ENVIRONMENT_ID=payload['EnvironmentID'], OAC_RUNTIME_SESSION_ID=payload['SessionID'], OAC_RUNTIME_NETWORK_ACCESS=payload['NetworkAccess'], diff --git a/services/agents-api/deploy/e2b/managed_init_test.py b/services/agents-api/deploy/e2b/managed_init_test.py index 8dadcaa6d..a166b79c4 100644 --- a/services/agents-api/deploy/e2b/managed_init_test.py +++ b/services/agents-api/deploy/e2b/managed_init_test.py @@ -41,6 +41,8 @@ def exercise(self, failed=False): 'OAC_RUNTIME_WORKSPACE': '/environment/workspace'} with patch.object(managed_init.shared, 'ROOT', root), patch.object(managed_init.shared, 'PROFILE', profile), \ patch.object(managed_init.shared, 'prepare_runtime', return_value=image_env), \ + patch.object(managed_init, 'SUSPEND_DIRECTORY', Path(temporary) / 'control'), \ + patch.object(managed_init.os, 'chown') as chown, \ patch.object(managed_init.os, 'fchown'), patch.object(managed_init.subprocess, 'Popen', process): if failed: with self.assertRaises(RuntimeError): @@ -58,6 +60,10 @@ def exercise(self, failed=False): self.assertNotIn(data['RuntimeBootstrap']['credential'], json.dumps(process.call_args.kwargs['env'])) self.assertEqual(process.call_args.kwargs['env']['OAC_RUNTIME_ENVIRONMENT_ID'], data['EnvironmentID']) self.assertEqual(process.call_args.kwargs['user'], 1000) + control = Path(temporary) / 'control' + self.assertEqual(control.stat().st_mode & 0o777, 0o700) + chown.assert_called_once_with(control, 1000, 1000) + self.assertEqual(process.call_args.kwargs['env']['OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE'], str(control / 'daemon-suspend.json')) if failed: self.assertFalse((root / 'managed-ready.json').exists()) else: diff --git a/services/agents-api/internal/execution/runtime_lifecycle.go b/services/agents-api/internal/execution/runtime_lifecycle.go index 8846a9bd3..655b6be5e 100644 --- a/services/agents-api/internal/execution/runtime_lifecycle.go +++ b/services/agents-api/internal/execution/runtime_lifecycle.go @@ -40,19 +40,20 @@ type RuntimeProvider struct { } type runtimeLifecycle struct { - store *store.Store - registry *gateway.Registry - config RuntimeProvider - nodeID string - gate chan struct{} - ctx context.Context - stop context.CancelFunc - cancelMu sync.Mutex - reconcileCancel context.CancelFunc - cursor string - pendingCursor string - connections map[string]*runtimeConnection - wakeHints chan struct{} + store *store.Store + registry *gateway.Registry + config RuntimeProvider + loadSuspensionPolicy func() *RuntimeSuspensionPolicy + nodeID string + gate chan struct{} + ctx context.Context + stop context.CancelFunc + cancelMu sync.Mutex + reconcileCancel context.CancelFunc + cursor string + pendingCursor string + connections map[string]*runtimeConnection + wakeHints chan struct{} } func newRuntimeManager(s *store.Store, registry *gateway.Registry, config *RuntimeProvider) (*runtimeManager, error) { @@ -119,6 +120,9 @@ func (r *runtimeLifecycle) lock(ctx context.Context) error { <-r.gate return err } + if r.loadSuspensionPolicy != nil { + r.config.Suspension = r.loadSuspensionPolicy() + } return nil case <-ctx.Done(): return ctx.Err() diff --git a/services/agents-api/internal/execution/runtime_manager.go b/services/agents-api/internal/execution/runtime_manager.go index 3678c408c..c0e3522b6 100644 --- a/services/agents-api/internal/execution/runtime_manager.go +++ b/services/agents-api/internal/execution/runtime_manager.go @@ -82,7 +82,8 @@ func (m *runtimeManager) node(id string) (*runtimeNode, error) { ctx, stop := context.WithCancel(m.ctx) n = &runtimeNode{lifecycle: &runtimeLifecycle{ store: m.store, registry: m.registry, config: m.config, nodeID: id, - gate: make(chan struct{}, 1), ctx: ctx, stop: stop, + loadSuspensionPolicy: m.suspensionPolicy, + gate: make(chan struct{}, 1), ctx: ctx, stop: stop, connections: make(map[string]*runtimeConnection), wakeHints: make(chan struct{}, 1), }} m.nodes[id] = n @@ -311,3 +312,15 @@ func (m *runtimeManager) stop() { // stop must precede drain. Both background loops and external provisioning or // manual reconciliation finish before Worker releases its unique writer lease. func (m *runtimeManager) drain() { m.active.Wait() } + +// Each serialized lifecycle operation uses the latest committed scheduling policy. +// Provider identity and generation routing remain bound to the existing lifecycle. +func (m *runtimeManager) suspensionPolicy() *RuntimeSuspensionPolicy { + m.mu.Lock() + defer m.mu.Unlock() + if m.config.Suspension == nil { + return nil + } + policy := *m.config.Suspension + return &policy +} diff --git a/services/agents-api/internal/execution/runtime_manager_test.go b/services/agents-api/internal/execution/runtime_manager_test.go index 9fc6b82ef..526d1854f 100644 --- a/services/agents-api/internal/execution/runtime_manager_test.go +++ b/services/agents-api/internal/execution/runtime_manager_test.go @@ -3,6 +3,7 @@ package execution import ( "context" "errors" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" "sync" "testing" "time" @@ -156,3 +157,45 @@ func TestRuntimeManagerHintsRemainPerNode(t *testing.T) { default: } } + +func TestRuntimeLifecycleLoadsCommittedPolicyAtOperationBoundary(t *testing.T) { + m := testRuntimeManager(t) + node, err := m.node("node") + if err != nil { + t.Fatal(err) + } + r := node.lifecycle + if err := r.lock(t.Context()); err != nil { + t.Fatal(err) + } + if r.config.Suspension != nil { + t.Fatal("unexpected initial policy") + } + policy := &RuntimeSuspensionPolicy{IdleTimeout: 5 * time.Minute, Retention: 24 * time.Hour, MaxActive: 4, MaxRetained: 16} + m.publishDeployment(PreparedRuntimeDeployment{Config: &RuntimeProvider{ProviderKind: "docker", Suspension: policy}}, store.RuntimeDeploymentView{Generation: 2}) + if r.config.Suspension != nil { + t.Fatal("in-flight operation policy changed") + } + <-r.gate + if err := r.lock(t.Context()); err != nil { + t.Fatal(err) + } + if r.config.Suspension == nil || *r.config.Suspension != *policy { + t.Fatal("existing lifecycle did not load committed policy") + } + if r.config.Suspension == policy { + t.Fatal("operation must own its policy snapshot") + } + <-r.gate + m.publishDeployment(PreparedRuntimeDeployment{Config: &RuntimeProvider{ProviderKind: "docker"}}, store.RuntimeDeploymentView{Generation: 3}) + if err := r.lock(t.Context()); err != nil { + t.Fatal(err) + } + defer func() { <-r.gate }() + if r.config.Suspension != nil { + t.Fatal("disabled policy was not adopted") + } + if again, err := m.node("node"); err != nil || again != node { + t.Fatal("policy update replaced lifecycle") + } +} From fd321da9cdc1c2465887e702606c86ac20615cbf Mon Sep 17 00:00:00 2001 From: sam Date: Wed, 30 Sep 2026 16:13:51 +0800 Subject: [PATCH 13/14] refactor(runtime): centralize hosted suspend control path --- apps/docs/content/docs/runtime-bootstrap.mdx | 13 ++++++++++ apps/docs/content/guide-sources.json | 4 +-- docs/runtime-bootstrap.md | 13 ++++++++++ internal/runtimebootstrap/suspension.go | 26 +++++++++++++++++++ internal/runtimebootstrap/suspension.json | 3 +++ .../agents-api/deploy/e2b/build-template.py | 2 ++ .../deploy/e2b/build_template_test.py | 3 +++ services/agents-api/deploy/e2b/init.py | 2 ++ services/agents-api/deploy/e2b/init_test.py | 4 ++- .../agents-api/deploy/e2b/managed_init.py | 14 +++++----- .../deploy/e2b/managed_init_test.py | 6 ++--- .../execution/runtime_compute_resident.go | 3 ++- .../execution/runtime_compute_wake.go | 3 ++- .../tools/microsandbox-provider/bootstrap.go | 5 ++-- 14 files changed, 84 insertions(+), 17 deletions(-) create mode 100644 internal/runtimebootstrap/suspension.go create mode 100644 internal/runtimebootstrap/suspension.json diff --git a/apps/docs/content/docs/runtime-bootstrap.mdx b/apps/docs/content/docs/runtime-bootstrap.mdx index 4dff42d41..9ace6402d 100644 --- a/apps/docs/content/docs/runtime-bootstrap.mdx +++ b/apps/docs/content/docs/runtime-bootstrap.mdx @@ -52,6 +52,19 @@ Self-hosted enrollment exchanges its executor credential for a daemon identity through the machine API; it is a different source of authority, not a managed bootstrap-file fallback. Both paths enter the same Runtime execution loop. +## Hosted suspension control + +The private control-file path is defined once in +[`suspension.json`](https://github.com/MiniMax-AI/parsar-core/blob/f6d258735fc601c521dd990e6f9e1ed261f4ef2d/internal/runtimebootstrap/suspension.json). Core recovery +commands and Go providers read it through `runtimebootstrap.SuspendControlFile()`. +The E2B template packager writes the same value into its protected Runtime +environment configuration; managed startup reads that value and prepares its +parent directory with mode 0700 and Runtime ownership. The daemon enables the +existing suspension protocol through `OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE`. +This is a packaged internal protocol setting, not an operator-editable idle policy +or a Provider lease timeout. The control file is not part of the bootstrap +credential document above. + ## Verification `go test ./internal/runtimebootstrap ./apps/parsar-daemon/internal/cli` covers the diff --git a/apps/docs/content/guide-sources.json b/apps/docs/content/guide-sources.json index 3727c3986..29c394ca0 100644 --- a/apps/docs/content/guide-sources.json +++ b/apps/docs/content/guide-sources.json @@ -26,7 +26,7 @@ "docs/assets/development-architecture.png": "24e6d0145d4f16ad70b07b6bc643808a6455aaf6398434d199cf74def200fca6", "docs/development.md": "f5c340253036a2cabc43e22aecdf70522d90280d6511e7649278ae93712ab8b6", "contracts/agents-api/harness-onboarding.md": "874a80dc1ffc5e97b2783bea3f6b217893f997b38bb8180a29f09c6dbf75e622", - "docs/runtime-bootstrap.md": "0d49aed73b298039e04453e6f465b0e925fb2227fa35d4206820bb3acbd8df39", + "docs/runtime-bootstrap.md": "f2ab66e6cb8a804e06eee39e4ca8723b0ae2a94e42d6b08aafc13d5bae52924e", "docs/runtime-protocol.md": "e8aa4cf862b5f63a4138cfeb25196a6bdb5c40a2e389e828b464585415dc6c45", "docs/sandbox-provider.md": "a20da12b5963461128e7a7581320a3eeddd6c0eb0dbdc10e35c3ab0a9e7c12fe", "apps/docs/scripts/guides.json": "3c3768fb94fd3d42464d4ce8c55724b9ba8360133c7cc19d513d8ec83a8e034f" @@ -57,7 +57,7 @@ "public/images/source/docs/assets/development-architecture.png": "24e6d0145d4f16ad70b07b6bc643808a6455aaf6398434d199cf74def200fca6", "content/docs/development.mdx": "2e5249ddfca571d264e93fd1e0e390a4bd5b0b623bda100cd02f2982929850bb", "content/docs/harness-onboarding.mdx": "17626e9f6256ddfffc95fd81dfa8d9c712d9ff16792ea5ef54bfe8c3ce1a2a9f", - "content/docs/runtime-bootstrap.mdx": "58982955811c9ad46a9fc04d8d2aef5a762fc9d25d5833f88aa75ee3fc46523f", + "content/docs/runtime-bootstrap.mdx": "7539eec21719ada6b73c31733f1e54b552cd219cebf589f7e917f4631dca0a1a", "content/docs/runtime-protocol.mdx": "d6b818edf2a0e0c4f2d3878f75805043c6dba3b09c6563f80faa4cecd05d40c7", "content/docs/sandbox-provider.mdx": "eeeab2d7963423928ca00e46c31e6a4d1ec56bb9ede442e3b6249f5114590fc2" } diff --git a/docs/runtime-bootstrap.md b/docs/runtime-bootstrap.md index c448f9b96..8a3bea830 100644 --- a/docs/runtime-bootstrap.md +++ b/docs/runtime-bootstrap.md @@ -49,6 +49,19 @@ Self-hosted enrollment exchanges its executor credential for a daemon identity through the machine API; it is a different source of authority, not a managed bootstrap-file fallback. Both paths enter the same Runtime execution loop. +## Hosted suspension control + +The private control-file path is defined once in +[`suspension.json`](../internal/runtimebootstrap/suspension.json). Core recovery +commands and Go providers read it through `runtimebootstrap.SuspendControlFile()`. +The E2B template packager writes the same value into its protected Runtime +environment configuration; managed startup reads that value and prepares its +parent directory with mode 0700 and Runtime ownership. The daemon enables the +existing suspension protocol through `OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE`. +This is a packaged internal protocol setting, not an operator-editable idle policy +or a Provider lease timeout. The control file is not part of the bootstrap +credential document above. + ## Verification `go test ./internal/runtimebootstrap ./apps/parsar-daemon/internal/cli` covers the diff --git a/internal/runtimebootstrap/suspension.go b/internal/runtimebootstrap/suspension.go new file mode 100644 index 000000000..68af5f37c --- /dev/null +++ b/internal/runtimebootstrap/suspension.go @@ -0,0 +1,26 @@ +package runtimebootstrap + +import ( + _ "embed" + "encoding/json" +) + +// suspensionJSON is also consumed by the E2B template packager. Keep hosted +// startup and recovery on the same private control-file location. +// +//go:embed suspension.json +var suspensionJSON []byte + +var suspendControlFile = func() string { + var configuration struct { + ControlFile string `json:"control_file"` + } + if err := json.Unmarshal(suspensionJSON, &configuration); err != nil { + panic(err) + } + return configuration.ControlFile +}() + +// SuspendControlFile is the packaged hosted Runtime protocol default, not a +// user-editable deployment setting. Providers prepare its private parent directory. +func SuspendControlFile() string { return suspendControlFile } diff --git a/internal/runtimebootstrap/suspension.json b/internal/runtimebootstrap/suspension.json new file mode 100644 index 000000000..682e8523a --- /dev/null +++ b/internal/runtimebootstrap/suspension.json @@ -0,0 +1,3 @@ +{ + "control_file": "/run/oac/daemon-suspend.json" +} diff --git a/services/agents-api/deploy/e2b/build-template.py b/services/agents-api/deploy/e2b/build-template.py index 9eb56cef4..4f99bcab4 100644 --- a/services/agents-api/deploy/e2b/build-template.py +++ b/services/agents-api/deploy/e2b/build-template.py @@ -27,6 +27,8 @@ if value.startswith(('HOME=', 'OAC_'))) if environment.get('OAC_RUNTIME_WORKSPACE') != '/environment/workspace': parser.error('Image does not use the colocated Runtime layout') +suspension = json.loads((Path(__file__).resolve().parents[4] / 'internal/runtimebootstrap/suspension.json').read_text()) +environment['OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE'] = suspension['control_file'] dev_home = Path(os.environ.get('OAC_DEV_HOME') or Path.home() / '.oac') if not dev_home.is_absolute(): parser.error('OAC_DEV_HOME must be absolute') diff --git a/services/agents-api/deploy/e2b/build_template_test.py b/services/agents-api/deploy/e2b/build_template_test.py index e44d89bb2..5bb447b69 100644 --- a/services/agents-api/deploy/e2b/build_template_test.py +++ b/services/agents-api/deploy/e2b/build_template_test.py @@ -64,6 +64,9 @@ def build(instance, **kwargs): self.assertIs(instance, template) context = Path(factory.call_args.kwargs['file_context_path']) self.assertEqual(stat.S_IMODE(context.stat().st_mode), 0o700) + defaults = json.loads((Path(__file__).resolve().parents[4] / 'internal/runtimebootstrap/suspension.json').read_text()) + runtime_env = json.loads((context / 'runtime-env.json').read_text()) + self.assertEqual(runtime_env['OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE'], defaults['control_file']) with tarfile.open(context / 'runtime.tar.gz') as archive: modes = {m.name: stat.S_IMODE(m.mode) for m in archive.getmembers()} for parent in ('usr', 'usr/local', 'etc'): diff --git a/services/agents-api/deploy/e2b/init.py b/services/agents-api/deploy/e2b/init.py index 75d8e1f07..45a590175 100644 --- a/services/agents-api/deploy/e2b/init.py +++ b/services/agents-api/deploy/e2b/init.py @@ -106,6 +106,8 @@ def initialize(): # Claim before any side effect. An interrupted attempt must never start twice. write_private(ROOT / 'launch.json', identity) environment = prepare_runtime() + # Application-owned startup does not participate in Core-managed suspension. + environment.pop("OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE", None) credential = PROFILE.parent / 'executor-key.json' write_private(credential, payload['executor_key'], owner=1000) source.unlink() diff --git a/services/agents-api/deploy/e2b/init_test.py b/services/agents-api/deploy/e2b/init_test.py index 3c416c13f..31a1489d3 100644 --- a/services/agents-api/deploy/e2b/init_test.py +++ b/services/agents-api/deploy/e2b/init_test.py @@ -37,7 +37,8 @@ def test_launch_handoff_is_private_and_cannot_replay(self): profile = Path(temporary) / 'private/default' environment_file = Path(temporary) / 'image.json' environment_file.write_text(json.dumps({'HOME': '/home/runtime', - 'OAC_RUNTIME_HOME': '/home/runtime/.oac', 'OAC_RUNTIME_WORKSPACE': '/environment/workspace'})) + 'OAC_RUNTIME_HOME': '/home/runtime/.oac', 'OAC_RUNTIME_WORKSPACE': '/environment/workspace', + 'OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE': '/private/control/fixture.json'})) (root / 'bootstrap.json').write_text(json.dumps(PAYLOAD)) real_chmod = Path.chmod @@ -57,6 +58,7 @@ def chmod(path, mode): self.assertEqual(argv[argv.index('--remote') + 1], PAYLOAD['remote_url']) self.assertNotIn('test-private-key', repr(popen.call_args)) self.assertNotIn('OAC_RUNTIME_SESSION_ID', options['env']) + self.assertNotIn('OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE', options['env']) self.assertEqual((options['user'], options['group'], options['extra_groups']), (1000, 1000, [])) self.assertEqual(options['umask'], 0o077) key = profile.parent / 'executor-key.json' diff --git a/services/agents-api/deploy/e2b/managed_init.py b/services/agents-api/deploy/e2b/managed_init.py index 615010dfb..a733c8541 100644 --- a/services/agents-api/deploy/e2b/managed_init.py +++ b/services/agents-api/deploy/e2b/managed_init.py @@ -13,9 +13,6 @@ _spec.loader.exec_module(shared) -SUSPEND_DIRECTORY = Path("/run/oac") - - def identity(payload): fields = ['InstallationID', 'TenantID', 'EnvironmentID', 'AllocationID', 'SessionID', 'DeviceID'] if not isinstance(payload, dict) or set(payload) != set(fields + ['NetworkAccess', 'AllowedDomains', 'RuntimeBootstrap']): @@ -45,10 +42,13 @@ def initialize(): binding = identity(payload) shared.write_private(root / 'managed-launch.json', binding) environment = shared.prepare_runtime() - SUSPEND_DIRECTORY.mkdir(mode=0o700, parents=True, exist_ok=True) - os.chown(SUSPEND_DIRECTORY, 1000, 1000) - SUSPEND_DIRECTORY.chmod(0o700) - environment['OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE'] = str(SUSPEND_DIRECTORY / 'daemon-suspend.json') + control_file = Path(environment['OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE']) + if not control_file.is_absolute() or '..' in control_file.parts: + raise ValueError('Absolute private suspend control file required') + control_directory = control_file.parent + control_directory.mkdir(mode=0o700, parents=True, exist_ok=True) + os.chown(control_directory, 1000, 1000) + control_directory.chmod(0o700) environment.update(OAC_RUNTIME_ENVIRONMENT_ID=payload['EnvironmentID'], OAC_RUNTIME_SESSION_ID=payload['SessionID'], OAC_RUNTIME_NETWORK_ACCESS=payload['NetworkAccess'], diff --git a/services/agents-api/deploy/e2b/managed_init_test.py b/services/agents-api/deploy/e2b/managed_init_test.py index a166b79c4..aa2a9a7b6 100644 --- a/services/agents-api/deploy/e2b/managed_init_test.py +++ b/services/agents-api/deploy/e2b/managed_init_test.py @@ -38,10 +38,10 @@ def exercise(self, failed=False): if failed: process.side_effect = RuntimeError('private process diagnostic') image_env = {'PATH': '/usr/local/bin:/usr/bin:/bin', 'OAC_RUNTIME_HOME': str(Path(temporary) / '.oac'), - 'OAC_RUNTIME_WORKSPACE': '/environment/workspace'} + 'OAC_RUNTIME_WORKSPACE': '/environment/workspace', + 'OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE': str(Path(temporary) / 'control' / 'custom-suspend.json')} with patch.object(managed_init.shared, 'ROOT', root), patch.object(managed_init.shared, 'PROFILE', profile), \ patch.object(managed_init.shared, 'prepare_runtime', return_value=image_env), \ - patch.object(managed_init, 'SUSPEND_DIRECTORY', Path(temporary) / 'control'), \ patch.object(managed_init.os, 'chown') as chown, \ patch.object(managed_init.os, 'fchown'), patch.object(managed_init.subprocess, 'Popen', process): if failed: @@ -63,7 +63,7 @@ def exercise(self, failed=False): control = Path(temporary) / 'control' self.assertEqual(control.stat().st_mode & 0o777, 0o700) chown.assert_called_once_with(control, 1000, 1000) - self.assertEqual(process.call_args.kwargs['env']['OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE'], str(control / 'daemon-suspend.json')) + self.assertEqual(process.call_args.kwargs['env']['OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE'], str(control / 'custom-suspend.json')) if failed: self.assertFalse((root / 'managed-ready.json').exists()) else: diff --git a/services/agents-api/internal/execution/runtime_compute_resident.go b/services/agents-api/internal/execution/runtime_compute_resident.go index fce5fdb82..bfcc5e5ac 100644 --- a/services/agents-api/internal/execution/runtime_compute_resident.go +++ b/services/agents-api/internal/execution/runtime_compute_resident.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "errors" + "github.com/MiniMax-AI-Dev/parsar/internal/runtimebootstrap" "time" "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/gateway" @@ -184,7 +185,7 @@ func (r *runtimeLifecycle) wakeResidentCompute(ctx context.Context, p sandbox.Re if !errors.Is(err, store.ErrNotFound) && !errors.Is(err, gateway.ErrDeviceNotRegistered) && !errors.Is(err, gateway.ErrSessionClosed) { return err } - result, err := p.RunCommand(ctx, runtimeReference(owner), sandbox.Command{Args: []string{"oac-daemon", "resume", "--control-file", "/run/oac/daemon-suspend.json", "--environment-id", owner.EnvironmentID, "--suspend-id", state.SuspendID}}) + result, err := p.RunCommand(ctx, runtimeReference(owner), sandbox.Command{Args: []string{"oac-daemon", "resume", "--control-file", runtimebootstrap.SuspendControlFile(), "--environment-id", owner.EnvironmentID, "--suspend-id", state.SuspendID}}) if err != nil { return err } diff --git a/services/agents-api/internal/execution/runtime_compute_wake.go b/services/agents-api/internal/execution/runtime_compute_wake.go index 5a1ad04ab..0decd6e00 100644 --- a/services/agents-api/internal/execution/runtime_compute_wake.go +++ b/services/agents-api/internal/execution/runtime_compute_wake.go @@ -3,6 +3,7 @@ package execution import ( "context" "errors" + "github.com/MiniMax-AI-Dev/parsar/internal/runtimebootstrap" "time" "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/gateway" @@ -24,7 +25,7 @@ func (r *runtimeLifecycle) wakeCompute(ctx context.Context, p sandbox.Checkpoint } // This idempotent control signal is fenced by guest PID/start time and the // suspension token. It cannot execute or replay an agent request. - result, err := p.RunCommandCompute(ctx, runtimeReference(owner), state.Current, sandbox.Command{Args: []string{"oac-daemon", "resume", "--control-file", "/run/oac/daemon-suspend.json", "--environment-id", owner.EnvironmentID, "--suspend-id", state.SuspendID}}) + result, err := p.RunCommandCompute(ctx, runtimeReference(owner), state.Current, sandbox.Command{Args: []string{"oac-daemon", "resume", "--control-file", runtimebootstrap.SuspendControlFile(), "--environment-id", owner.EnvironmentID, "--suspend-id", state.SuspendID}}) if err != nil { return err } diff --git a/services/agents-api/tools/microsandbox-provider/bootstrap.go b/services/agents-api/tools/microsandbox-provider/bootstrap.go index 7eef9c4c3..654c6b315 100644 --- a/services/agents-api/tools/microsandbox-provider/bootstrap.go +++ b/services/agents-api/tools/microsandbox-provider/bootstrap.go @@ -5,6 +5,7 @@ package main import ( "context" "encoding/json" + "github.com/MiniMax-AI-Dev/parsar/internal/runtimebootstrap" "time" "github.com/MiniMax-AI-Dev/parsar/internal/agentnetwork" @@ -18,7 +19,7 @@ import ( const bootstrapScript = ` import ctypes,json,os,stat,subprocess,sys b=json.load(sys.stdin) -for p in ['/home/runtime','/home/runtime/.oac','/environment','/environment/workspace','/environment/staging','/environment/initialization','/environment/packages','/run/oac']: +for p in ['/home/runtime','/home/runtime/.oac','/environment','/environment/workspace','/environment/staging','/environment/initialization','/environment/packages',os.path.dirname(os.environ['OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE'])]: os.makedirs(p,mode=0o700,exist_ok=True) if not stat.S_ISDIR(os.lstat(p).st_mode): raise RuntimeError('invalid bootstrap directory') os.chmod(p,0o700);os.chown(p,1000,1000) @@ -64,7 +65,7 @@ func (b backend) create(ctx context.Context) (wire.Response, error) { "HOME": "/home/runtime", "OAC_RUNTIME_HOME": "/home/runtime/.oac", "OAC_RUNTIME_ENVIRONMENT_ID": bootstrap.EnvironmentID, "OAC_RUNTIME_SESSION_ID": bootstrap.SessionID, "OAC_RUNTIME_NETWORK_ACCESS": policy.Access, "OAC_RUNTIME_ALLOWED_DOMAINS": string(domains), - "OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE": "/run/oac/daemon-suspend.json", + "OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE": runtimebootstrap.SuspendControlFile(), })) if e != nil { return wire.Response{}, e From 4860acd44f5caccaa7d8ec6fc0b6cf0ec90e47e8 Mon Sep 17 00:00:00 2001 From: sam Date: Wed, 30 Sep 2026 17:05:50 +0800 Subject: [PATCH 14/14] Validate resident suspension policy in provider registration --- docs/sandbox-provider.md | 2 +- .../internal/sandbox/providers/registration.go | 16 ++++++++++------ .../sandbox/providers/registration_test.go | 14 +++++++++++++- 3 files changed, 24 insertions(+), 8 deletions(-) diff --git a/docs/sandbox-provider.md b/docs/sandbox-provider.md index ab9eeda7d..9ccdc60ee 100644 --- a/docs/sandbox-provider.md +++ b/docs/sandbox-provider.md @@ -229,7 +229,7 @@ runtime plugin loading. - A `nodes` registration requires only `BuildLocal`; a `direct` registration requires only `BuildDirect`. Missing, mixed or unknown modes are rejected. - Specification/resource validators, the configuration adapter and the complete operation declaration are mandatory. An incomplete registration cannot publish a partial installer projection. - The Runtime input policy must either accept the pinned Runtime or give the adapter's fixed reason for rejecting it; it cannot do both. -- Current checkpoint suspension requires node mode and positive idle/retention defaults that fit Runtime durations. The two durations are independent. Providers without checkpoint support must not configure suspension defaults. +- Checkpoint suspension requires node mode; resident pause requires direct mode. Either lifecycle requires positive idle/retention defaults that fit Runtime durations. The two durations are independent. Providers supporting neither lifecycle must not configure suspension defaults. The configuration adapter must be non-nil, including its concrete value. Each `ConfigurationRequirements` field needs an explicit valid decision: `Credential` and `PublicOrigin` use `Required` or `NotRequired`; `Discovery` uses the shared supported/Unsupported declaration with its safe reason. `ConfigurationDiscoverer` must be implemented even when discovery is unsupported. New requirement fields or discovery methods require an explicit validation update; they cannot inherit an existing decision. Configuration discovery is distinct from resource selection discovery. Requiring a credential does not itself promise the resource operation `VerifyCredential`. diff --git a/services/core/internal/sandbox/providers/registration.go b/services/core/internal/sandbox/providers/registration.go index 66ea85dde..5f700fcbb 100644 --- a/services/core/internal/sandbox/providers/registration.go +++ b/services/core/internal/sandbox/providers/registration.go @@ -45,16 +45,20 @@ func ValidateRegistration(a Adapter) error { if err := sandbox.ValidateOperations(operations); err != nil { return err } - // The current common lifecycle admits checkpoint suspension only on nodes, - // and creates its policy whenever checkpoint support is declared. - if operations["Initial"].State == providercontract.Supported { + // Both suspension lifecycles use the common idle/retention policy. + checkpoint := operations["Initial"].State == providercontract.Supported + resident := operations["PauseResident"].State == providercontract.Supported + if (checkpoint && a.Mode != "nodes") || (resident && a.Mode != "direct") { + return invalid("suspension mode") + } + if checkpoint || resident { const maximumSeconds = int64((1<<63 - 1) / time.Second) - if a.Mode != "nodes" || a.IdleSeconds < 1 || a.RetentionSeconds < 1 || + if a.IdleSeconds < 1 || a.RetentionSeconds < 1 || a.IdleSeconds > maximumSeconds || a.RetentionSeconds > maximumSeconds { - return invalid("checkpoint policy") + return invalid("suspension policy") } } else if a.IdleSeconds != 0 || a.RetentionSeconds != 0 { - return invalid("non-checkpoint policy") + return invalid("unsupported suspension policy") } return nil } diff --git a/services/core/internal/sandbox/providers/registration_test.go b/services/core/internal/sandbox/providers/registration_test.go index 764eee16a..a6590ee96 100644 --- a/services/core/internal/sandbox/providers/registration_test.go +++ b/services/core/internal/sandbox/providers/registration_test.go @@ -158,13 +158,17 @@ func TestCompleteRegistrationsPreserveConstruction(t *testing.T) { // Idle time is measured before suspension, retention after suspension. Neither // duration needs to be greater than the other. -func TestRegistrationCheckpointPolicy(t *testing.T) { +func TestRegistrationSuspensionPolicy(t *testing.T) { for _, tc := range []struct { name string kind string idle, retention int64 direct, valid bool }{ + {"resident suspension", "e2b", 300, 86400, false, true}, + {"resident missing idle", "e2b", 0, 86400, false, false}, + {"resident missing retention", "e2b", 300, 0, false, false}, + {"resident overflow", "e2b", 1<<63 - 1, 86400, false, false}, {"negative idle", "microsandbox", -1, 20, false, false}, {"missing idle", "microsandbox", 0, 20, false, false}, {"missing retention", "microsandbox", 20, 0, false, false}, @@ -187,3 +191,11 @@ func TestRegistrationCheckpointPolicy(t *testing.T) { }) } } + +func TestRegistrationRejectsResidentPauseOnNodeTransport(t *testing.T) { + a := adapters["e2b"] + a.Mode, a.BuildDirect, a.BuildLocal = "nodes", nil, adapters["docker"].BuildLocal + if err := ValidateRegistration(a); !errors.Is(err, providercontract.ErrContract) { + t.Fatal(err) + } +}