diff --git a/contracts/agents-api/sandbox-deployment.md b/contracts/agents-api/sandbox-deployment.md index 360c27fd3..fd786c89b 100644 --- a/contracts/agents-api/sandbox-deployment.md +++ b/contracts/agents-api/sandbox-deployment.md @@ -84,12 +84,13 @@ For E2B, post `{"configuration": {"api_url": "…", "domain": "…"}, "credentia GET and successful writes return `installation_id`, `provider`, `core_url` (read-only: the installation public URL, present even before configuration), `mode`, `generation`, `owner_epoch`, `reset`, `rollout`, `suspension`, `resources` and `credential_configured`. A configured deployment also returns `specification`, `specification_digest`, `configuration` and `metadata`: the adapter's public projection of its selectors and of the observations it recorded, never raw stored values or secrets. E2B returns `configuration.template`, `configuration.api_url`, `configuration.domain` and, once recorded, `metadata.template_build`. Docker and microsandbox return empty `configuration` and `metadata` objects and `credential_configured: false`; an unconfigured deployment has neither object. - `metadata.template_build` is `{status, resources: {cpus, memory_mib, root_disk_mib}}`: the build as Core read it through the pinned SDK when the selection was saved. GET never calls E2B, so it stays cheap during an E2B outage. Validation admits only a `ready` build whose CPU count and memory equal the selected values; `root_disk_mib` is the build's native disk size, which Core does not enforce. Unknown values are null, `metadata: {}` means no observation was recorded, and an identical PUT without a credential does not refresh it. -- `suspension` is `{idle_seconds, retention_seconds}` for microsandbox, the only provider Core suspends (currently 300 and 86400); Docker, E2B and unconfigured deployments return null. - Request `resources` and response `specification.resources` are per-sandbox limits. Response `resources.allocations` and `resources.pending` count unreleased allocations and pending hosted Environments without an allocation. - An unconfigured deployment has an empty provider and no specification. Docker and microsandbox use `mode: nodes`; E2B uses `mode: direct`, without a synthetic node. - `generation` identifies the saved selection. `owner_epoch` fences the execution owner and node connections; it does not replace `expected_generation`. - `specification_digest` is the server's identity of the provider, limits and Runtime release; enrollment echoes it unchanged. +`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. + The typed `SandboxAdminClient` in `packages/agents-client` checks the deployment, node list, node detail and allocation responses against exactly these shapes. An unknown or missing member, or a wrong type, rejects the whole response with a 502 `invalid_admin_response` error. A node's `diagnostic` is absent or a code, never empty, and the client reads an unknown code as `provider_unavailable`. ## Initial setup and same-provider changes @@ -190,7 +191,7 @@ Some fields keep one name across providers but differ in meaning, or do not appl | 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; `configuration.template` selects the build | 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 `metadata.template_build` | The build as Core read it when the selection was saved | Absent: `metadata` is empty | Absent: `metadata` is empty | -| 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 | @@ -200,7 +201,7 @@ Some fields keep one name across providers but differ in meaning, or do not appl | 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` | -| Runtime observation `lifecycle_state: sleeping` | Never | Never | While suspended | +| Runtime observation `lifecycle_state: sleeping` | No running sample while paused | 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 | ## Errors diff --git a/docs/runtime-bootstrap.md b/docs/runtime-bootstrap.md index 9a1ca5e96..0041b16c1 100644 --- a/docs/runtime-bootstrap.md +++ b/docs/runtime-bootstrap.md @@ -29,6 +29,19 @@ The Runtime validates the input and owns authentication and connection. A succes Self-hosted executors and operator-provisioned devices get their daemon identity in other ways; the [machine connection API](../contracts/agents-api/machine-api.md#credentials) lists every credential source. All of them 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/daemon/internal/cli` covers the input contract, the exclusivity of credential sources and restart behavior. Provider tests verify delivery and file permissions without relying on the Runtime's private storage. diff --git a/docs/sandbox-provider.md b/docs/sandbox-provider.md index 1cf4b37d8..c0b49e3e1 100644 --- a/docs/sandbox-provider.md +++ b/docs/sandbox-provider.md @@ -11,6 +11,8 @@ A **Sandbox Provider** supplies the outer compute that a Runtime daemon runs in Core owns durable Environment, allocation, placement and cleanup state; the Provider owns compute and bootstrap only. The Runtime prepares capabilities and runs Turns over the [Core–Runtime protocol](runtime-protocol.md), and the provider hands it its identity through the [Runtime bootstrap](runtime-bootstrap.md) file. A provider never runs Environment initialization, Skills, Plugins, MCP setup, initial files, execution or Files; those use the Runtime. Isolation belongs to the provider's infrastructure, not the daemon; see [Runtime and outer isolation](concepts.md#runtime-and-outer-isolation). Use the vendor's maintained SDK behind a thin adapter. +A thin adapter lets a Provider be replaced without changing the common execution flow. Core owns shared scheduling, persistence and recovery through declared capability contracts; provider SDK calls and native behavior stay inside the adapter. Shared policy decisions use registered policy rather than vendor names. + ## Steps 1. **Read the contract.** Implement the five required operations and give an explicit decision for every extension interface in [Implement the interface](#implement-the-interface). @@ -47,9 +49,12 @@ Every provider implements the methods of each interface below and returns a comp | `runtimeobs.BatchSource` | Explicit decision; requires `Source` | Bounded observations in input order, with per-target errors | | `sandbox.SelectionDiscoverer` | Explicit decision | Read-only native configuration discovery before commit | | `sandbox.CredentialVerifier` | Explicit decision | Verify access to owned resources without mutation | +| `sandbox.ResidentPauseProvider` | Explicit decision for both operations | Memory-preserving pause and resume of the same owned compute ID | `CheckpointProvider` adds `Initial` and `NewCompute` (construct compute references without allocating), `GetCompute`, `Suspend`, `Resume`, `ResumeCompute` (thaw only the same resident instance after an aborted pause), `KillCompute`, `DeleteSnapshot` and `RunCommandCompute`, which runs a bounded command in one exact compute incarnation. Core uses `RunCommandCompute` to wake a parked daemon after a restore ([`runtime_compute_wake.go`](../services/core/internal/execution/runtime_compute_wake.go)). +The checkpoint lifecycle and resident pause/resume pair must each agree on support. + Each declaration entry is `state: supported` with no reason, or `state: unsupported` with an authored reason code. Missing, zero, unknown or unsafe entries and missing methods fail validation. Adding a method to an interface requires an explicit decision and implementation in every adapter; never supply a base type or generate blanket unsupported implementations. An unsupported method returns `providercontract.UnsupportedError` before any native I/O. The error names the exact operation and a safe code, never a native message, resource identity, endpoint or credential. An empty result, a nil error, `Unavailable` or an unknown mutation outcome never stands in for unsupported, and the five required methods can never return it. @@ -90,6 +95,12 @@ Every call receives a bounded context. Expiry or cancellation ends the caller's Checkpoint support adds `Compute` generation, name and ID and `SnapshotIdentity`; persist operation IDs and the provider's snapshot provenance unchanged. `ObserveOnly` on suspend or resume observes the previous attempt and never starts another capture or restore. `ResumeCompute` only thaws the retained source and never cold-starts a stopped one. Cleanup targets the exact compute incarnation and snapshot, not whatever instance now has the same name. Read [`runtime_compute.go`](../services/core/internal/execution/runtime_compute.go) and its failure tests before declaring checkpoint support. +Resident pause retains the original provider ID and has no separate snapshot +identity. Core quiesces the Runtime before calling `PauseResident`, persists the phase +before the native call, and admits work only after `ResumeResident` and Runtime +reconnection. An unknown pause result must be observed without replaying a +late mutation. See [the resident lifecycle](../services/core/internal/execution/runtime_compute_resident.go). + ### Four distinct readiness facts | Fact | Evidence | Does not establish | @@ -124,7 +135,7 @@ A new provider takes these steps: - A `nodes` registration has only `BuildLocal`, and a `direct` registration only `BuildDirect`; missing, mixed or unknown modes are rejected. - The specification and resource validators, the configuration adapter and a complete operation declaration are mandatory, so an incomplete registration cannot publish a partial installer projection. - The Runtime input policy either accepts the pinned Runtime or gives the adapter's fixed reason for rejecting it, never both. -- Checkpoint support requires node mode and positive idle and retention defaults that fit Runtime durations; a provider without checkpoint support configures no 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. Every `ConfigurationRequirements` field needs an explicit valid decision: `Credential` and `PublicOrigin` are `Required` or `NotRequired`, and `Discovery` uses the shared supported or unsupported declaration with a safe reason. A new requirement field or discovery method needs an explicit validation update and never inherits an existing decision. Configuration discovery is distinct from resource selection discovery, and requiring a credential does not promise the `VerifyCredential` operation. These checks establish complete registration, not correct native SDK behavior; constructor and adapter contract tests still apply. @@ -134,7 +145,7 @@ Preview and persistence use `providers.Normalize` and `providers.Describe`. `Sel A direct adapter with a credential verifies all retained generations and allocation references before a key is replaced. The common `sandbox.CallFence` excludes native calls and waits for helper completion, including calls whose callers timed out; execution invokes the prepared verification and fencing callbacks without branching on a vendor. -Vendor deployment validation and SDK setup stay at the construction boundary, and construction never creates an Environment. For node-local adapters `providers.Built` returns the provider, probe, installation identity, backend fingerprint and specification digest, and the factory also returns its close function. `execution.RuntimeProvider` binds the adapter to its kind, installation ID, backend fingerprint, generation, mode and node ownership; the database owns the selection, and the in-memory copy is never another authority. Docker and microsandbox run on nodes, and E2B is constructed directly. The node proxy exposes checkpoint operations only for a backend whose registered declaration supports them, and common lifecycle code admits suspension through `CheckpointProvider`, never through a provider name. +Vendor deployment validation and SDK setup stay at the construction boundary, and construction never creates an Environment. For node-local adapters `providers.Built` returns the provider, probe, installation identity, backend fingerprint and specification digest, and the factory also returns its close function. `execution.RuntimeProvider` binds the adapter to its kind, installation ID, backend fingerprint, generation, mode and node ownership; the database owns the selection, and the in-memory copy is never another authority. Docker and microsandbox run on nodes, and E2B is constructed directly. The node proxy exposes checkpoint operations only for a backend whose registered declaration supports them, and common lifecycle code admits suspension through the declared `CheckpointProvider` or `ResidentPauseProvider` capability, never through a provider name. The backend fingerprint identifies a native resource namespace, not capacity. Core keeps deployment generations so that owned allocations keep resolving to their original backend; never repoint retained allocations at a replacement backend. @@ -182,7 +193,7 @@ The deployment's CPU, memory and disk settings, `max_active`, `max_retained` and ### Suspension -A provider with checkpoint support can suspend idle work; the deployment's [`suspension`](../contracts/agents-api/sandbox-deployment.md#safe-response) policy sets the idle time and snapshot retention. Core suspends only after at least one Turn is terminal, when no root or Subagent Turn is queued, in progress or waiting, no input, file operation or initialization is pending, and real activity has been idle for the configured interval. For node allocations Core records the first root or child terminal transition with the database clock in the same transaction. Candidate filtering and the Session-locked recheck compare elapsed database time with the idle duration, and the initial snapshot retention deadline is anchored to the same database observation, so Core and database host clocks need not agree. Native completion timestamps stay unchanged in public history but never drive idle admission, and heartbeats never reset activity. Before acknowledging a planned suspension, the daemon closes admission and drains native cleanup, output receipts and file work. +A provider with checkpoint or resident pause support can suspend idle work; the deployment's [`suspension`](../contracts/agents-api/sandbox-deployment.md#safe-response) policy sets the idle time and snapshot retention. Checkpoint suspension requires at least one terminal Turn; resident pause also admits initialized Sessions that have never received a Turn. Core suspends when no root or Subagent Turn is queued, in progress or waiting, no input, file operation or initialization is pending, and real activity has been idle for the configured interval. For node allocations Core records the first root or child terminal transition with the database clock in the same transaction. Candidate filtering and the Session-locked recheck compare elapsed database time with the idle duration, and the initial snapshot retention deadline is anchored to the same database observation, so Core and database host clocks need not agree. Native completion timestamps stay unchanged in public history but never drive idle admission, and heartbeats never reset activity. Before acknowledging a planned suspension, the daemon closes admission and drains native cleanup, output receipts and file work. The Worker lease, the Session lock and the per-node gates own suspension for every provider. New Turn claims, file-write intents and capture admission serialize under the Session lock and share one compute-phase check; new pending work cancels a capture and wakes the same source. Normal preparation waits for the compute phase to be running, after the authenticated resume handshake, and pending input stays pending when its promotion conflicts with a lifecycle transition. Compute phases and revision-checked receipts live on the allocation. Core persists quiesce, capture and restore intent before the effect, only a fresh receipt performs a capture or restore, and recovery observes the exact attempt without retrying an unknown creation, capture or restore. A consumed snapshot never rolls a running generation back. Deletion, revocation and retention expiry win over wake, up to the final database compare-and-swap, and unknown cleanup identities are kept until owned resources are confirmed absent. Consumed artifacts and old compute are deleted, so suspension cycles never build a chain of writable disks. 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/core/cmd/server/managed_generations_resident.go b/services/core/cmd/server/managed_generations_resident.go new file mode 100644 index 000000000..757554b5f --- /dev/null +++ b/services/core/cmd/server/managed_generations_resident.go @@ -0,0 +1,44 @@ +package main + +import ( + "context" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/providercontract" + + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" +) + +// Resident operations retain the allocation generation and credential fence. + +func (p *generationRouter) PauseResident(ctx context.Context, ref sandbox.Reference) (sandbox.Info, error) { + if err := providercontract.Require(p, "PauseResident"); err != nil { + return sandbox.Info{}, err + } + provider, done, err := p.route(ctx, ref) + if err != nil { + return sandbox.Info{}, err + } + defer done() + resident, err := sandbox.ResidentPause(provider) + if err != nil { + return sandbox.Info{}, err + } + return resident.PauseResident(ctx, ref) +} + +func (p *generationRouter) ResumeResident(ctx context.Context, ref sandbox.Reference) (sandbox.Info, error) { + if err := providercontract.Require(p, "ResumeResident"); err != nil { + return sandbox.Info{}, err + } + provider, done, err := p.route(ctx, ref) + if err != nil { + return sandbox.Info{}, err + } + defer done() + resident, err := sandbox.ResidentPause(provider) + if err != nil { + return sandbox.Info{}, err + } + return resident.ResumeResident(ctx, ref) +} + +var _ sandbox.ResidentPauseProvider = (*generationRouter)(nil) diff --git a/services/core/cmd/server/managed_generations_test.go b/services/core/cmd/server/managed_generations_test.go index aca3167eb..225637045 100644 --- a/services/core/cmd/server/managed_generations_test.go +++ b/services/core/cmd/server/managed_generations_test.go @@ -10,6 +10,8 @@ import ( "testing" "time" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox/e2b" "github.com/MiniMax-AI/OpenAgentCore/services/core/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 { @@ -60,7 +63,22 @@ print(json.dumps({'Version':1,'Info':info})) db := &routingSetupStore{setupStore: setupStore{value: current}, old: old, oldID: ref.AllocationID} setup := &managedSetup{processPaths: paths, 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, capabilityErr := sandbox.ResidentPause(router) + if capabilityErr != nil { + 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 { @@ -69,6 +87,12 @@ print(json.dumps({'Version':1,'Info':info})) if _, err := router.Renew(ctx, ref); err != nil { t.Fatal(err) } + if _, err := resident.PauseResident(ctx, ref); err != nil { + t.Fatal(err) + } + if _, err := resident.ResumeResident(ctx, ref); err != nil { + t.Fatal(err) + } if err := router.Kill(ctx, ref); err != nil { t.Fatal(err) } @@ -82,7 +106,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 { @@ -96,7 +120,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.Configuration.(*e2b.DeploymentConfiguration).Template || q.Config.Resources != expected.Specification.Resources { diff --git a/services/core/cmd/server/managed_setup.go b/services/core/cmd/server/managed_setup.go index 3ec771c93..34acd8c5a 100644 --- a/services/core/cmd/server/managed_setup.go +++ b/services/core/cmd/server/managed_setup.go @@ -126,7 +126,7 @@ 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 sandbox.SupportsCheckpoint(provider) { + if (sandbox.SupportsCheckpoint(provider) || sandbox.SupportsResidentPause(provider)) && 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/core/deploy/e2b/build-template.py b/services/core/deploy/e2b/build-template.py index d78926213..f124bf95a 100644 --- a/services/core/deploy/e2b/build-template.py +++ b/services/core/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/core/deploy/e2b/build_template_test.py b/services/core/deploy/e2b/build_template_test.py index 9a3b3034f..2886f1c3c 100644 --- a/services/core/deploy/e2b/build_template_test.py +++ b/services/core/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/core/deploy/e2b/helper_contract_generated.py b/services/core/deploy/e2b/helper_contract_generated.py index ebd7fadf2..987dc1e48 100644 --- a/services/core/deploy/e2b/helper_contract_generated.py +++ b/services/core/deploy/e2b/helper_contract_generated.py @@ -10,7 +10,7 @@ MAX_REQUEST = 75497472 MAX_RESPONSE = 16777216 NETWORK_ACCESS = ["enabled","disabled","restricted"] -OPERATIONS = ["create","inspect","renew","kill","command","validate_deployment","observe","list_templates","list_builds","verify_credential"] +OPERATIONS = ["create","inspect","renew","pause","resume","kill","command","validate_deployment","observe","list_templates","list_builds","verify_credential"] PROTOCOL_VERSION = 1 REFERENCE_FIELDS = ["TenantID","EnvironmentID","AllocationID"] REQUEST_FIELDS = ["Version","Operation","Config","Reference","References","Bootstrap","RuntimeBootstrap","Command","Deadline"] diff --git a/services/core/deploy/e2b/init.py b/services/core/deploy/e2b/init.py index 75d8e1f07..45a590175 100644 --- a/services/core/deploy/e2b/init.py +++ b/services/core/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/core/deploy/e2b/init_test.py b/services/core/deploy/e2b/init_test.py index 3c416c13f..31a1489d3 100644 --- a/services/core/deploy/e2b/init_test.py +++ b/services/core/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/core/deploy/e2b/managed_init.py b/services/core/deploy/e2b/managed_init.py index a64931f14..a332e08ce 100644 --- a/services/core/deploy/e2b/managed_init.py +++ b/services/core/deploy/e2b/managed_init.py @@ -47,6 +47,13 @@ def initialize(): binding = identity(payload) shared.write_private(root / 'managed-launch.json', binding) environment = shared.prepare_runtime() + 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/core/deploy/e2b/managed_init_test.py b/services/core/deploy/e2b/managed_init_test.py index d9017ba82..53eb3c9ad 100644 --- a/services/core/deploy/e2b/managed_init_test.py +++ b/services/core/deploy/e2b/managed_init_test.py @@ -63,9 +63,11 @@ 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.os, 'chown') as chown, \ patch.object(managed_init.os, 'fchown'), patch.object(managed_init.subprocess, 'Popen', process): if failed: with self.assertRaises(RuntimeError): @@ -83,6 +85,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 / 'custom-suspend.json')) if failed: self.assertFalse((root / 'managed-ready.json').exists()) else: diff --git a/services/core/internal/api/sandbox_deployment_changes_test.go b/services/core/internal/api/sandbox_deployment_changes_test.go index e270e7a9a..83f866458 100644 --- a/services/core/internal/api/sandbox_deployment_changes_test.go +++ b/services/core/internal/api/sandbox_deployment_changes_test.go @@ -10,6 +10,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox/e2b" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/device" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/adminaudit" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store" ) @@ -17,8 +18,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.Configuration == nil || in.Configuration.(*e2b.DeploymentConfiguration).APIKey != "synthetic-private-key" { t.Fatal("write-only fields were lost") } diff --git a/services/core/internal/api/sandbox_deployment_setup.go b/services/core/internal/api/sandbox_deployment_setup.go index bd68951a9..796147fd9 100644 --- a/services/core/internal/api/sandbox_deployment_setup.go +++ b/services/core/internal/api/sandbox_deployment_setup.go @@ -114,6 +114,7 @@ func (h *Handler) updateSandboxDeployment(w http.ResponseWriter, r *http.Request writeStoreError(w, r, err) return } + setAdminAuditSource(r, "") result, err := h.sandboxUpdate(r.Context(), store.SandboxDeploymentUpdateRequest{SandboxDeploymentSetupRequest: selection, ExpectedGeneration: *input.ExpectedGeneration}) if err != nil { writeStoreError(w, r, err) diff --git a/services/core/internal/execution/provider_operations_fixture_test.go b/services/core/internal/execution/provider_operations_fixture_test.go index 4b568704b..0c78e2669 100644 --- a/services/core/internal/execution/provider_operations_fixture_test.go +++ b/services/core/internal/execution/provider_operations_fixture_test.go @@ -9,6 +9,8 @@ import ( func (*lifecycleOnlySandbox) ProviderOperations() providercontract.Operations { return providercontract.Operations{ + "PauseResident": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "ResumeResident": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, "Create": {State: providercontract.Supported}, "GetInfo": {State: providercontract.Supported}, "Renew": {State: providercontract.Supported}, @@ -75,3 +77,11 @@ func (*lifecycleOnlySandbox) ObservationProviderType() string { return "fixture" func (p *lifecycleOnlySandbox) ResolveObservationSource(context.Context) (runtimeobs.Source, error) { return p, nil } + +func (*lifecycleOnlySandbox) PauseResident(context.Context, sandbox.Reference) (sandbox.Info, error) { + return sandbox.Info{}, &providercontract.UnsupportedError{Operation: "PauseResident", Reason: "fixture_operation_not_supported"} +} + +func (*lifecycleOnlySandbox) ResumeResident(context.Context, sandbox.Reference) (sandbox.Info, error) { + return sandbox.Info{}, &providercontract.UnsupportedError{Operation: "ResumeResident", Reason: "fixture_operation_not_supported"} +} diff --git a/services/core/internal/execution/runtime_compute.go b/services/core/internal/execution/runtime_compute.go index 055d0e2ba..2aeb44280 100644 --- a/services/core/internal/execution/runtime_compute.go +++ b/services/core/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,30 @@ func (r *runtimeLifecycle) saveCompute(ctx context.Context, owner store.RuntimeA } idleTimeout = policy.IdleTimeout } + resident := sandbox.SupportsResidentPause(r.config.Provider) + 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 sandbox.SupportsResidentPause(r.config.Provider) { + p, err := sandbox.ResidentPause(r.config.Provider) + if err != nil { + return err + } + 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, capabilityErr := sandbox.Checkpoint(r.config.Provider) if capabilityErr != nil { return capabilityErr @@ -83,6 +104,13 @@ func (r *runtimeLifecycle) enableCompute(ctx context.Context, owner store.Runtim } func (r *runtimeLifecycle) observeCompute(ctx context.Context, owner store.RuntimeAllocation) error { + if sandbox.SupportsResidentPause(r.config.Provider) { + p, err := sandbox.ResidentPause(r.config.Provider) + if err != nil { + return err + } + return r.observeResidentCompute(ctx, p, owner) + } p, capabilityErr := sandbox.Checkpoint(r.config.Provider) if capabilityErr != nil { return capabilityErr diff --git a/services/core/internal/execution/runtime_compute_resident.go b/services/core/internal/execution/runtime_compute_resident.go new file mode 100644 index 000000000..78ab1b97b --- /dev/null +++ b/services/core/internal/execution/runtime_compute_resident.go @@ -0,0 +1,226 @@ +package execution + +import ( + "context" + "encoding/json" + "errors" + "github.com/MiniMax-AI/OpenAgentCore/internal/runtimebootstrap" + "time" + + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/gateway" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" + "github.com/MiniMax-AI/OpenAgentCore/services/core/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": + // 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 + } + 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 + } + 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 + } + if _, err = r.store.KeepRuntimeAllocation(ctx, owner); err != nil { + return err + } + if activity.WakeRequested { + return r.store.ClearRuntimeWake(ctx, owner, owner.ComputeActivityAt) + } + return nil + } + 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.PauseResident(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.ResumeResident(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", runtimebootstrap.SuspendControlFile(), "--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/core/internal/execution/runtime_compute_wake.go b/services/core/internal/execution/runtime_compute_wake.go index 7b01fefd8..5c04cb375 100644 --- a/services/core/internal/execution/runtime_compute_wake.go +++ b/services/core/internal/execution/runtime_compute_wake.go @@ -3,6 +3,7 @@ package execution import ( "context" "errors" + "github.com/MiniMax-AI/OpenAgentCore/internal/runtimebootstrap" "time" "github.com/MiniMax-AI/OpenAgentCore/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/core/internal/execution/runtime_lifecycle.go b/services/core/internal/execution/runtime_lifecycle.go index 0f0aefabf..561d0f343 100644 --- a/services/core/internal/execution/runtime_lifecycle.go +++ b/services/core/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) { @@ -96,12 +97,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 !sandbox.SupportsCheckpoint(config.Provider) || policy.IdleTimeout < time.Second || policy.Retention < time.Second || policy.MaxActive < 1 || policy.MaxRetained < policy.MaxActive { + checkpoint := sandbox.SupportsCheckpoint(config.Provider) + resident := sandbox.SupportsResidentPause(config.Provider) + 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 @@ -120,6 +123,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/core/internal/execution/runtime_manager.go b/services/core/internal/execution/runtime_manager.go index 68fb2f4b3..e9980276e 100644 --- a/services/core/internal/execution/runtime_manager.go +++ b/services/core/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/core/internal/execution/runtime_manager_test.go b/services/core/internal/execution/runtime_manager_test.go index 9fc6b82ef..034f3f928 100644 --- a/services/core/internal/execution/runtime_manager_test.go +++ b/services/core/internal/execution/runtime_manager_test.go @@ -3,6 +3,7 @@ package execution import ( "context" "errors" + "github.com/MiniMax-AI/OpenAgentCore/services/core/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") + } +} diff --git a/services/core/internal/sandbox/docker/operations.go b/services/core/internal/sandbox/docker/operations.go index c2758008d..7fbe41d28 100644 --- a/services/core/internal/sandbox/docker/operations.go +++ b/services/core/internal/sandbox/docker/operations.go @@ -15,6 +15,8 @@ var _ runtimeobs.BatchSource = (*Provider)(nil) // Operations is this adapter's complete authored resource contract. func Operations() providercontract.Operations { return providercontract.Operations{ + "PauseResident": {State: providercontract.Unsupported, Reason: "docker_does_not_support_resident_pause"}, + "ResumeResident": {State: providercontract.Unsupported, Reason: "docker_does_not_support_resident_pause"}, "Create": {State: providercontract.Supported}, "GetInfo": {State: providercontract.Supported}, "Renew": {State: providercontract.Supported}, @@ -74,3 +76,11 @@ func (p *Provider) DiscoverSelection(context.Context, sandbox.Selection) (sandbo func (p *Provider) VerifyCredential(context.Context, []sandbox.Reference) error { return &providercontract.UnsupportedError{Operation: "VerifyCredential", Reason: Operations()["VerifyCredential"].Reason} } + +func (*Provider) PauseResident(context.Context, sandbox.Reference) (sandbox.Info, error) { + return sandbox.Info{}, &providercontract.UnsupportedError{Operation: "PauseResident", Reason: "docker_does_not_support_resident_pause"} +} + +func (*Provider) ResumeResident(context.Context, sandbox.Reference) (sandbox.Info, error) { + return sandbox.Info{}, &providercontract.UnsupportedError{Operation: "ResumeResident", Reason: "docker_does_not_support_resident_pause"} +} diff --git a/services/core/internal/sandbox/e2b/helper_contract.go b/services/core/internal/sandbox/e2b/helper_contract.go index 7e9a96f25..480df5c51 100644 --- a/services/core/internal/sandbox/e2b/helper_contract.go +++ b/services/core/internal/sandbox/e2b/helper_contract.go @@ -22,7 +22,7 @@ const MaxCommandInputBytes = sandbox.MaxCommandInputBytes // HelperOperations declares the complete set of one-shot helper operations. func HelperOperations() []string { - return []string{"create", "inspect", "renew", "kill", "command", "validate_deployment", "observe", "list_templates", "list_builds", "verify_credential"} + return []string{"create", "inspect", "renew", "pause", "resume", "kill", "command", "validate_deployment", "observe", "list_templates", "list_builds", "verify_credential"} } // HelperErrors are sanitized wire outcomes; an empty code denotes success. diff --git a/services/core/internal/sandbox/e2b/operations.go b/services/core/internal/sandbox/e2b/operations.go index 14530ebf9..b96bc018c 100644 --- a/services/core/internal/sandbox/e2b/operations.go +++ b/services/core/internal/sandbox/e2b/operations.go @@ -15,6 +15,8 @@ var _ runtimeobs.BatchSource = (*Provider)(nil) // Operations is this adapter's complete authored resource contract. func Operations() providercontract.Operations { return providercontract.Operations{ + "PauseResident": {State: providercontract.Supported}, + "ResumeResident": {State: providercontract.Supported}, "Create": {State: providercontract.Supported}, "GetInfo": {State: providercontract.Supported}, "Renew": {State: providercontract.Supported}, diff --git a/services/core/internal/sandbox/e2b/provider.go b/services/core/internal/sandbox/e2b/provider.go index c368d8866..680c27e64 100644 --- a/services/core/internal/sandbox/e2b/provider.go +++ b/services/core/internal/sandbox/e2b/provider.go @@ -108,6 +108,7 @@ type Provider struct { } var _ sandbox.SandboxProvider = (*Provider)(nil) +var _ sandbox.ResidentPauseProvider = (*Provider)(nil) func validID(value string) bool { id, err := uuid.Parse(value) @@ -252,6 +253,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) PauseResident(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) ResumeResident(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/core/internal/sandbox/microsandbox/operations.go b/services/core/internal/sandbox/microsandbox/operations.go index f59c9ed02..6efdf5f85 100644 --- a/services/core/internal/sandbox/microsandbox/operations.go +++ b/services/core/internal/sandbox/microsandbox/operations.go @@ -15,6 +15,8 @@ var _ runtimeobs.BatchSource = (*Provider)(nil) // Operations is this adapter's complete authored resource contract. func Operations() providercontract.Operations { return providercontract.Operations{ + "PauseResident": {State: providercontract.Unsupported, Reason: "microsandbox_does_not_support_resident_pause"}, + "ResumeResident": {State: providercontract.Unsupported, Reason: "microsandbox_does_not_support_resident_pause"}, "Create": {State: providercontract.Supported}, "GetInfo": {State: providercontract.Supported}, "Renew": {State: providercontract.Supported}, @@ -47,3 +49,11 @@ func (p *Provider) DiscoverSelection(context.Context, sandbox.Selection) (sandbo func (p *Provider) VerifyCredential(context.Context, []sandbox.Reference) error { return &providercontract.UnsupportedError{Operation: "VerifyCredential", Reason: Operations()["VerifyCredential"].Reason} } + +func (*Provider) PauseResident(context.Context, sandbox.Reference) (sandbox.Info, error) { + return sandbox.Info{}, &providercontract.UnsupportedError{Operation: "PauseResident", Reason: "microsandbox_does_not_support_resident_pause"} +} + +func (*Provider) ResumeResident(context.Context, sandbox.Reference) (sandbox.Info, error) { + return sandbox.Info{}, &providercontract.UnsupportedError{Operation: "ResumeResident", Reason: "microsandbox_does_not_support_resident_pause"} +} diff --git a/services/core/internal/sandbox/node/provider_operations_fixture_test.go b/services/core/internal/sandbox/node/provider_operations_fixture_test.go index 8d34c1511..c666a4a78 100644 --- a/services/core/internal/sandbox/node/provider_operations_fixture_test.go +++ b/services/core/internal/sandbox/node/provider_operations_fixture_test.go @@ -9,6 +9,8 @@ import ( func (*fakeProvider) ProviderOperations() providercontract.Operations { return providercontract.Operations{ + "PauseResident": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "ResumeResident": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, "Create": {State: providercontract.Supported}, "GetInfo": {State: providercontract.Supported}, "Renew": {State: providercontract.Supported}, @@ -72,6 +74,8 @@ func (*fakeProvider) VerifyCredential(context.Context, []sandbox.Reference) erro } func (*observationProvider) ProviderOperations() providercontract.Operations { return providercontract.Operations{ + "PauseResident": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "ResumeResident": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, "Create": {State: providercontract.Supported}, "GetInfo": {State: providercontract.Supported}, "Renew": {State: providercontract.Supported}, @@ -104,3 +108,11 @@ func (*observationProvider) ObservationProviderType() string { return "fixture" func (p *observationProvider) ResolveObservationSource(context.Context) (runtimeobs.Source, error) { return p, nil } + +func (*fakeProvider) PauseResident(context.Context, sandbox.Reference) (sandbox.Info, error) { + return sandbox.Info{}, &providercontract.UnsupportedError{Operation: "PauseResident", Reason: "fixture_operation_not_supported"} +} + +func (*fakeProvider) ResumeResident(context.Context, sandbox.Reference) (sandbox.Info, error) { + return sandbox.Info{}, &providercontract.UnsupportedError{Operation: "ResumeResident", Reason: "fixture_operation_not_supported"} +} diff --git a/services/core/internal/sandbox/node/proxy.go b/services/core/internal/sandbox/node/proxy.go index cf30880ea..72068acf0 100644 --- a/services/core/internal/sandbox/node/proxy.go +++ b/services/core/internal/sandbox/node/proxy.go @@ -75,6 +75,8 @@ func (p *provider) ProviderOperations() providercontract.Operations { return nil } ops := adapter.Operations() + ops["PauseResident"] = providercontract.Support{State: providercontract.Unsupported, Reason: "node_transport_has_no_resident_pause"} + ops["ResumeResident"] = providercontract.Support{State: providercontract.Unsupported, Reason: "node_transport_has_no_resident_pause"} ops["DiscoverSelection"] = providercontract.Support{State: providercontract.Unsupported, Reason: "node_configuration_is_core_owned"} ops["VerifyCredential"] = providercontract.Support{State: providercontract.Unsupported, Reason: "node_credentials_are_transport_owned"} ops["ObserveBatch"] = providercontract.Support{State: providercontract.Unsupported, Reason: "node_transport_has_no_batch_observation"} @@ -175,3 +177,11 @@ func (h *Hub) GenerationProvider(kind string, resolve func(context.Context, sand p := &provider{hub: h, kind: kind, resolveGeneration: resolve} return p } + +func (*provider) PauseResident(context.Context, sandbox.Reference) (sandbox.Info, error) { + return sandbox.Info{}, &providercontract.UnsupportedError{Operation: "PauseResident", Reason: "node_transport_has_no_resident_pause"} +} + +func (*provider) ResumeResident(context.Context, sandbox.Reference) (sandbox.Info, error) { + return sandbox.Info{}, &providercontract.UnsupportedError{Operation: "ResumeResident", Reason: "node_transport_has_no_resident_pause"} +} diff --git a/services/core/internal/sandbox/operations.go b/services/core/internal/sandbox/operations.go index 3faf4454b..9d6c194ea 100644 --- a/services/core/internal/sandbox/operations.go +++ b/services/core/internal/sandbox/operations.go @@ -11,7 +11,7 @@ import ( // These existing interfaces are the canonical operation inventory. Declarations // must cover every method, including explicit unsupported implementations. var providerInterfaces = []reflect.Type{ - reflect.TypeFor[SandboxProvider](), reflect.TypeFor[CheckpointProvider](), + reflect.TypeFor[SandboxProvider](), reflect.TypeFor[CheckpointProvider](), reflect.TypeFor[ResidentPauseProvider](), reflect.TypeFor[SelectionDiscoverer](), reflect.TypeFor[CredentialVerifier](), reflect.TypeFor[runtimeobs.SourceResolver](), reflect.TypeFor[runtimeobs.Source](), reflect.TypeFor[runtimeobs.BatchSource](), } @@ -66,6 +66,9 @@ func ValidateOperations(operations providercontract.Operations) error { return fmt.Errorf("%w: incomplete checkpoint lifecycle", providercontract.ErrContract) } } + if operations["PauseResident"].State != operations["ResumeResident"].State { + return fmt.Errorf("%w: incomplete resident lifecycle", providercontract.ErrContract) + } if operations["ObserveBatch"].State == providercontract.Supported && operations["Observe"].State != providercontract.Supported { return fmt.Errorf("%w: batch observation requires observation", providercontract.ErrContract) } @@ -86,3 +89,17 @@ func Checkpoint(p SandboxProvider) (CheckpointProvider, error) { } return cp, nil } + +func SupportsResidentPause(p SandboxProvider) bool { + return providercontract.Require(p, "PauseResident") == nil +} +func ResidentPause(p SandboxProvider) (ResidentPauseProvider, error) { + if err := providercontract.Require(p, "PauseResident"); err != nil { + return nil, err + } + resident, ok := p.(ResidentPauseProvider) + if !ok { + return nil, providercontract.ErrContract + } + return resident, nil +} diff --git a/services/core/internal/sandbox/operations_test.go b/services/core/internal/sandbox/operations_test.go index f1cc091d1..ea47822c0 100644 --- a/services/core/internal/sandbox/operations_test.go +++ b/services/core/internal/sandbox/operations_test.go @@ -125,3 +125,21 @@ func TestProviderRegistrationRequiresObservationIdentity(t *testing.T) { t.Fatal("provider registration accepted empty observation identity", err) } } + +func TestResidentSupportUsesDeclarationRatherThanMethodPresence(t *testing.T) { + for _, p := range []sandbox.SandboxProvider{&docker.Provider{}, µsandbox.Provider{}} { + if sandbox.SupportsResidentPause(p) { + t.Fatalf("%T advertises resident pause", p) + } + if _, err := sandbox.ResidentPause(p); !errors.Is(err, providercontract.ErrUnsupported) { + t.Fatalf("%T: %v", p, err) + } + } + p := &e2b.Provider{} + if !sandbox.SupportsResidentPause(p) { + t.Fatal("E2B resident pause missing") + } + if _, err := sandbox.ResidentPause(p); err != nil { + t.Fatal(err) + } +} diff --git a/services/core/internal/sandbox/providers/registration.go b/services/core/internal/sandbox/providers/registration.go index 712283e55..228eb40c5 100644 --- a/services/core/internal/sandbox/providers/registration.go +++ b/services/core/internal/sandbox/providers/registration.go @@ -51,16 +51,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 65fc3a355..b053b1ba6 100644 --- a/services/core/internal/sandbox/providers/registration_test.go +++ b/services/core/internal/sandbox/providers/registration_test.go @@ -168,13 +168,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}, @@ -197,3 +201,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) + } +} diff --git a/services/core/internal/sandbox/providers/registry.go b/services/core/internal/sandbox/providers/registry.go index f579441f7..4f29bba0b 100644 --- a/services/core/internal/sandbox/providers/registry.go +++ b/services/core/internal/sandbox/providers/registry.go @@ -50,6 +50,7 @@ var adapters = map[string]Adapter{ }, "e2b": { Policy: e2b.Policy(), Operations: e2b.Operations, Mode: "direct", BuildDirect: buildE2B, + IdleSeconds: 300, RetentionSeconds: 86400, Configuration: e2b.ConfigurationAdapter{}, ValidateSpecification: e2b.ValidateSpecification, ValidateResources: e2b.ValidateResources, }, diff --git a/services/core/internal/sandbox/providers/registry_test.go b/services/core/internal/sandbox/providers/registry_test.go index 3d3639b13..d40eeca59 100644 --- a/services/core/internal/sandbox/providers/registry_test.go +++ b/services/core/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/core/internal/sandbox/sandbox_provider.go b/services/core/internal/sandbox/sandbox_provider.go index 663526b8e..4af91bd99 100644 --- a/services/core/internal/sandbox/sandbox_provider.go +++ b/services/core/internal/sandbox/sandbox_provider.go @@ -92,6 +92,16 @@ type SandboxProvider interface { RunCommand(context.Context, Reference, Command) (CommandResult, error) } +// ResidentPauseProvider keeps the same owned compute incarnation across a +// memory-preserving pause. PauseResident must never create compute, and ResumeResident must +// reconnect only the existing allocation. Both operations must reconcile an +// unknown result against provider state before reporting success. +type ResidentPauseProvider interface { + SandboxProvider + PauseResident(context.Context, Reference) (Info, error) + ResumeResident(context.Context, Reference) (Info, error) +} + // CheckpointProvider is an explicitly declared extension. It supplies // exact-incarnation operations; Worker and Store remain the lifecycle owner. type CheckpointProvider interface { diff --git a/services/core/internal/store/archived_cancellation_migration_test.go b/services/core/internal/store/archived_cancellation_migration_test.go index 0ffa52a2d..e5ec38788 100644 --- a/services/core/internal/store/archived_cancellation_migration_test.go +++ b/services/core/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/core/internal/store/e2b_idle_downgrade_test.go b/services/core/internal/store/e2b_idle_downgrade_test.go new file mode 100644 index 000000000..693d9e4cf --- /dev/null +++ b/services/core/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/core/internal/store/provider_operations_fixture_test.go b/services/core/internal/store/provider_operations_fixture_test.go index 68634162e..7e7e1ebb1 100644 --- a/services/core/internal/store/provider_operations_fixture_test.go +++ b/services/core/internal/store/provider_operations_fixture_test.go @@ -9,6 +9,8 @@ import ( func (*lifecycleProvider) ProviderOperations() providercontract.Operations { return providercontract.Operations{ + "PauseResident": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "ResumeResident": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, "Create": {State: providercontract.Supported}, "GetInfo": {State: providercontract.Supported}, "Renew": {State: providercontract.Supported}, @@ -72,6 +74,8 @@ func (*lifecycleProvider) VerifyCredential(context.Context, []sandbox.Reference) } func (*fakeCheckpointProvider) ProviderOperations() providercontract.Operations { return providercontract.Operations{ + "PauseResident": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "ResumeResident": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, "Create": {State: providercontract.Supported}, "GetInfo": {State: providercontract.Supported}, "Renew": {State: providercontract.Supported}, @@ -104,3 +108,11 @@ func (*fakeCheckpointProvider) ObservationProviderType() string { return "fixtur func (p *fakeCheckpointProvider) ResolveObservationSource(context.Context) (runtimeobs.Source, error) { return p, nil } + +func (*lifecycleProvider) PauseResident(context.Context, sandbox.Reference) (sandbox.Info, error) { + return sandbox.Info{}, &providercontract.UnsupportedError{Operation: "PauseResident", Reason: "fixture_operation_not_supported"} +} + +func (*lifecycleProvider) ResumeResident(context.Context, sandbox.Reference) (sandbox.Info, error) { + return sandbox.Info{}, &providercontract.UnsupportedError{Operation: "ResumeResident", Reason: "fixture_operation_not_supported"} +} diff --git a/services/core/internal/store/provider_registration_migration_test.go b/services/core/internal/store/provider_registration_migration_test.go index 06a6074c8..b2086f757 100644 --- a/services/core/internal/store/provider_registration_migration_test.go +++ b/services/core/internal/store/provider_registration_migration_test.go @@ -29,6 +29,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/core/internal/store/runtime_compute_lifecycle_test.go b/services/core/internal/store/runtime_compute_lifecycle_test.go index e7934adbc..0596216c2 100644 --- a/services/core/internal/store/runtime_compute_lifecycle_test.go +++ b/services/core/internal/store/runtime_compute_lifecycle_test.go @@ -18,6 +18,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto/prototest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox/e2b" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store" "github.com/google/uuid" "github.com/gorilla/websocket" @@ -261,6 +262,7 @@ type computeLifecycleFixture struct { store *store.Store pool *pgxpool.Pool provider *fakeCheckpointProvider + managed sandbox.SandboxProvider worker *execution.Worker stop func() key string @@ -268,6 +270,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() @@ -283,14 +289,34 @@ 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 sandbox.SupportsResidentPause(f.managed) { + 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 { t.Fatal(err) } @@ -300,6 +326,15 @@ func (f *computeLifecycleFixture) start() { } f.worker, f.stop = w, stop t.Cleanup(stop) + if sandbox.SupportsResidentPause(f.managed) { + _, err := w.InitializeSandboxDeployment(t.Context(), store.SandboxDeploymentSetupRequest{ + Provider: "e2b", DeploymentSpec: store.SandboxDeploymentTestSpec("e2b"), + Configuration: &e2b.DeploymentConfiguration{APIKey: "fixture-api-key", Template: "runtime:" + uuid.NewString(), TemplateBuild: &e2b.DeploymentBuild{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() @@ -313,13 +348,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 @@ -332,11 +367,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/core/internal/store/runtime_idle_policy_test.go b/services/core/internal/store/runtime_idle_policy_test.go index b08b740cf..2922240e6 100644 --- a/services/core/internal/store/runtime_idle_policy_test.go +++ b/services/core/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/core/internal/store/runtime_lifecycle_test.go b/services/core/internal/store/runtime_lifecycle_test.go index 9f489fa18..f51b64aa0 100644 --- a/services/core/internal/store/runtime_lifecycle_test.go +++ b/services/core/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 } diff --git a/services/core/internal/store/runtime_resident_lifecycle_test.go b/services/core/internal/store/runtime_resident_lifecycle_test.go new file mode 100644 index 000000000..386238520 --- /dev/null +++ b/services/core/internal/store/runtime_resident_lifecycle_test.go @@ -0,0 +1,200 @@ +package store_test + +import ( + "context" + "encoding/json" + "errors" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/providercontract" + "sync/atomic" + "testing" + + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" +) + +type fakeResidentProvider struct { + *fakeCheckpointProvider + pauses, resumes int + losePause bool + renewals atomic.Int32 +} + +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) PauseResident(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) ResumeResident(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) + 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 { + 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") + } +} + +func (p *fakeResidentProvider) ProviderOperations() providercontract.Operations { + ops := p.lifecycleProvider.ProviderOperations() + ops["PauseResident"] = providercontract.Support{State: providercontract.Supported} + ops["ResumeResident"] = providercontract.Support{State: providercontract.Supported} + return ops +} + +func (p *fakeResidentProvider) Renew(ctx context.Context, r sandbox.Reference) (sandbox.Info, error) { + p.renewals.Add(1) + return p.fakeCheckpointProvider.Renew(ctx, r) +} + +func TestResidentRepeatedWakeRequestsStillRenewRunningCompute(t *testing.T) { + var resident *fakeResidentProvider + f := newComputeLifecycleFixtureWithProvider(t, 2, 4, func(p *fakeCheckpointProvider) sandbox.SandboxProvider { + resident = &fakeResidentProvider{fakeCheckpointProvider: p} + return resident + }) + tenant, _, environment, _ := f.create() + for range 3 { + if err := f.store.TouchRuntimeActivity(t.Context(), tenant, environment.ID); err != nil { + t.Fatal(err) + } + before := resident.renewals.Load() + owner := f.phase(tenant, environment.ID, "running") + if resident.renewals.Load() <= before { + t.Fatal("wake request skipped provider renewal") + } + var wakeRequested bool + err := f.pool.QueryRow(t.Context(), `SELECT compute_wake_requested FROM runtime_allocations WHERE id=$1`, owner.ID).Scan(&wakeRequested) + if err != nil || wakeRequested { + t.Fatal("wake request not cleared after renewal", err) + } + } + if resident.pauses != 0 || resident.resumes != 0 { + t.Fatal("active file access changed resident compute") + } +} diff --git a/services/core/internal/store/runtime_suspension.go b/services/core/internal/store/runtime_suspension.go index 4661292ef..f780a1bda 100644 --- a/services/core/internal/store/runtime_suspension.go +++ b/services/core/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/core/internal/store/sandbox_deployment_mutations.go b/services/core/internal/store/sandbox_deployment_mutations.go index c53def61a..c013c3055 100644 --- a/services/core/internal/store/sandbox_deployment_mutations.go +++ b/services/core/internal/store/sandbox_deployment_mutations.go @@ -30,6 +30,15 @@ func (s *Store) sandboxSelectionEqual(d sqlc.RuntimeDeployment, input SandboxDep if d.ProviderKind != input.Provider { 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 + } previous, err := s.sandboxSetup(d) if err != nil { return false, err diff --git a/services/core/internal/store/sandbox_deployment_setup.go b/services/core/internal/store/sandbox_deployment_setup.go index fc636b06b..dea970f71 100644 --- a/services/core/internal/store/sandbox_deployment_setup.go +++ b/services/core/internal/store/sandbox_deployment_setup.go @@ -186,8 +186,7 @@ func runtimeDeploymentView(d sqlc.RuntimeDeployment, publicURL string) (RuntimeD result.Metadata = configurationJSON(record.Metadata) result.CredentialConfigured = len(d.ProviderCredential) > 0 } - - if providers.SupportsCheckpoint(d.ProviderKind) { + if d.IdleSeconds > 0 && d.RetentionSeconds > 0 { result.Suspension = &SandboxSuspensionView{IdleSeconds: d.IdleSeconds, RetentionSeconds: d.RetentionSeconds} } return result, nil diff --git a/services/core/internal/store/sandbox_deployment_view_test.go b/services/core/internal/store/sandbox_deployment_view_test.go index 21da06627..63ae8930c 100644 --- a/services/core/internal/store/sandbox_deployment_view_test.go +++ b/services/core/internal/store/sandbox_deployment_view_test.go @@ -42,7 +42,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) @@ -75,8 +75,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 || string(view.Configuration) != "{}" || view.Suspension == nil || view.Suspension.IdleSeconds != 300 || view.Suspension.RetentionSeconds != 86400 { t.Fatalf("microsandbox suspension view = %+v %v", view, err) diff --git a/services/core/internal/store/sandbox_generations_test.go b/services/core/internal/store/sandbox_generations_test.go index f3a0e8714..7aa62f755 100644 --- a/services/core/internal/store/sandbox_generations_test.go +++ b/services/core/internal/store/sandbox_generations_test.go @@ -290,6 +290,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) diff --git a/services/core/internal/store/sandbox_policy_comparison_test.go b/services/core/internal/store/sandbox_policy_comparison_test.go new file mode 100644 index 000000000..ed55fb392 --- /dev/null +++ b/services/core/internal/store/sandbox_policy_comparison_test.go @@ -0,0 +1,50 @@ +package store + +import ( + "encoding/json" + "testing" + + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/db/sqlc" + "github.com/MiniMax-AI/OpenAgentCore/services/core/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{WebManaged: true, 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") + } + }) + } + }) + } +} diff --git a/services/core/tools/e2b-provider/README.md b/services/core/tools/e2b-provider/README.md index aff992d1b..42fec44a9 100644 --- a/services/core/tools/e2b-provider/README.md +++ b/services/core/tools/e2b-provider/README.md @@ -15,6 +15,8 @@ The account key is stored encrypted in Core's database and is write-only. It rea | `create` | `Create` | Creates the sandbox once and runs `managed_init.py`; see [Create](#create) | | `inspect` | `GetInfo` | Reads the sandbox by recorded ID, or by ownership metadata when no ID is recorded, and checks ownership, domain, template and resources | | `renew` | `Renew` | Extends the lease of the running sandbox to the configured timeout, then rereads it | +| `pause` | `PauseResident` | Preserves memory of the same sandbox; persists the pending operation before calling the SDK and settles an unknown result through observation | +| `resume` | `ResumeResident` | Reconnects only the same sandbox ID with `on_resume="restore"` | | `kill` | `Kill` | Destroys every matching sandbox and confirms that none remains | | `command` | `RunCommand` | Runs one bounded command as the Runtime user on a running sandbox whose bootstrap completed; output is limited to 1 MiB per stream | | `validate_deployment` | Deployment setup | Reads the template's builds and requires the exact build to be ready with the configured CPU and memory. Without configured resources the selection adopts the build's CPU and memory. Returns the build's status, CPU, memory and reported disk for Core to record; bounded to 30 seconds | diff --git a/services/core/tools/e2b-provider/helper_contract_generated.py b/services/core/tools/e2b-provider/helper_contract_generated.py index ebd7fadf2..987dc1e48 100644 --- a/services/core/tools/e2b-provider/helper_contract_generated.py +++ b/services/core/tools/e2b-provider/helper_contract_generated.py @@ -10,7 +10,7 @@ MAX_REQUEST = 75497472 MAX_RESPONSE = 16777216 NETWORK_ACCESS = ["enabled","disabled","restricted"] -OPERATIONS = ["create","inspect","renew","kill","command","validate_deployment","observe","list_templates","list_builds","verify_credential"] +OPERATIONS = ["create","inspect","renew","pause","resume","kill","command","validate_deployment","observe","list_templates","list_builds","verify_credential"] PROTOCOL_VERSION = 1 REFERENCE_FIELDS = ["TenantID","EnvironmentID","AllocationID"] REQUEST_FIELDS = ["Version","Operation","Config","Reference","References","Bootstrap","RuntimeBootstrap","Command","Deadline"] diff --git a/services/core/tools/e2b-provider/provider.py b/services/core/tools/e2b-provider/provider.py index dbf3111f0..42c57c920 100644 --- a/services/core/tools/e2b-provider/provider.py +++ b/services/core/tools/e2b-provider/provider.py @@ -189,6 +189,9 @@ def inspect(self): cloud = found[0] self.qualified(cloud) record = self.receipt.data or {} + # A qualified observation settles a lost pause reply without replaying it. + if cloud.state == 'paused' and record.get('pause_status') == 'pending' and record.get('bootstrap_complete'): + self.receipt.save(pause_status='settled') if (cloud.state == 'running' and not record.get('bootstrap_complete') and record.get('connection') and record.get('status') not in ('bootstrap_failed', 'killed')): try: @@ -263,6 +266,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 @@ -419,7 +467,8 @@ def execute(self): raise Failure('unconfirmed') result = run(self.client(cloud), self.q['Command'], self.remaining) return {'Version': PROTOCOL_VERSION, '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': PROTOCOL_VERSION, 'Info': self.info(cloud, absent=cloud is None), 'ErrorCode': ''} except Failure as error: info = self.info() diff --git a/services/core/tools/e2b-provider/provider_test.py b/services/core/tools/e2b-provider/provider_test.py index 74d708547..45e084683 100644 --- a/services/core/tools/e2b-provider/provider_test.py +++ b/services/core/tools/e2b-provider/provider_test.py @@ -83,6 +83,46 @@ 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('inspect')['Info']['State'], 'paused') + self.assertEqual(self.record()['pause_status'], 'settled') + def connect(sandbox_id, **kwargs): + self.assertEqual(sandbox_id, self.cloud.sandbox_id) + self.cloud.state = 'running' + return self.cloud + self.api.connect.side_effect = connect + self.assertEqual(self.call('resume')['Info']['State'], 'running') + self.api.pause.assert_called_once() + self.api.connect.assert_called_once() + 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, diff --git a/services/core/tools/e2b-provider/testdata/contract.json b/services/core/tools/e2b-provider/testdata/contract.json index aa79da367..83369500c 100644 --- a/services/core/tools/e2b-provider/testdata/contract.json +++ b/services/core/tools/e2b-provider/testdata/contract.json @@ -37,9 +37,11 @@ {"kind":"request","name":"observe-references-type","payload":{"Bootstrap":{"AllocationID":"33333333-3333-4333-8333-333333333333","AllowedDomains":null,"CoreURL":"https://core.example/api/v1","Credential":"fixture-only","DeviceID":"55555555-5555-4555-8555-555555555555","EnvironmentID":"22222222-2222-4222-8222-222222222222","NetworkAccess":"enabled","SessionID":"44444444-4444-4444-8444-444444444444","TenantID":"11111111-1111-4111-8111-111111111111"},"Config":{"APIKey":"","APIURL":"","Binary":"","Domain":"","InstallationID":"66666666-6666-4666-8666-666666666666","Resources":null,"StateDir":"","Template":"","TimeoutSeconds":0},"Deadline":"2099-01-01T00:00:00Z","Operation":"observe","Reference":{"AllocationID":"33333333-3333-4333-8333-333333333333","EnvironmentID":"22222222-2222-4222-8222-222222222222","TenantID":"11111111-1111-4111-8111-111111111111"},"References":[null],"RuntimeBootstrap":{"core_url":"https://core.example/api/v1","credential":"fixture-only","device_id":"55555555-5555-4555-8555-555555555555","version":1},"Version":1},"valid":false}, {"kind":"request","name":"operation-null","payload":{"Bootstrap":{"AllocationID":"33333333-3333-4333-8333-333333333333","AllowedDomains":null,"CoreURL":"https://core.example/api/v1","Credential":"fixture-only","DeviceID":"55555555-5555-4555-8555-555555555555","EnvironmentID":"22222222-2222-4222-8222-222222222222","NetworkAccess":"enabled","SessionID":"44444444-4444-4444-8444-444444444444","TenantID":"11111111-1111-4111-8111-111111111111"},"Config":{"APIKey":"","APIURL":"","Binary":"","Domain":"","InstallationID":"66666666-6666-4666-8666-666666666666","Resources":null,"StateDir":"","Template":"","TimeoutSeconds":0},"Deadline":"2099-01-01T00:00:00Z","Operation":null,"Reference":{"AllocationID":"33333333-3333-4333-8333-333333333333","EnvironmentID":"22222222-2222-4222-8222-222222222222","TenantID":"11111111-1111-4111-8111-111111111111"},"RuntimeBootstrap":{"core_url":"https://core.example/api/v1","credential":"fixture-only","device_id":"55555555-5555-4555-8555-555555555555","version":1},"Version":1},"valid":false}, {"kind":"request","name":"operation-unknown","payload":{"Bootstrap":{"AllocationID":"33333333-3333-4333-8333-333333333333","AllowedDomains":null,"CoreURL":"https://core.example/api/v1","Credential":"fixture-only","DeviceID":"55555555-5555-4555-8555-555555555555","EnvironmentID":"22222222-2222-4222-8222-222222222222","NetworkAccess":"enabled","SessionID":"44444444-4444-4444-8444-444444444444","TenantID":"11111111-1111-4111-8111-111111111111"},"Config":{"APIKey":"","APIURL":"","Binary":"","Domain":"","InstallationID":"66666666-6666-4666-8666-666666666666","Resources":null,"StateDir":"","Template":"","TimeoutSeconds":0},"Deadline":"2099-01-01T00:00:00Z","Operation":"unsupported","Reference":{"AllocationID":"33333333-3333-4333-8333-333333333333","EnvironmentID":"22222222-2222-4222-8222-222222222222","TenantID":"11111111-1111-4111-8111-111111111111"},"RuntimeBootstrap":{"core_url":"https://core.example/api/v1","credential":"fixture-only","device_id":"55555555-5555-4555-8555-555555555555","version":1},"Version":1},"valid":false}, +{"kind":"request","name":"pause","payload":{"Version":1,"Operation":"pause","Config":{"Binary":"","StateDir":"","InstallationID":"66666666-6666-4666-8666-666666666666","APIKey":"","Template":"","APIURL":"","Domain":"","TimeoutSeconds":0,"Resources":null},"Reference":{"TenantID":"11111111-1111-4111-8111-111111111111","EnvironmentID":"22222222-2222-4222-8222-222222222222","AllocationID":"33333333-3333-4333-8333-333333333333"},"Bootstrap":{"TenantID":"11111111-1111-4111-8111-111111111111","EnvironmentID":"22222222-2222-4222-8222-222222222222","AllocationID":"33333333-3333-4333-8333-333333333333","SessionID":"44444444-4444-4444-8444-444444444444","DeviceID":"55555555-5555-4555-8555-555555555555","CoreURL":"https://core.example/api/v1","Credential":"fixture-only","NetworkAccess":"enabled","AllowedDomains":null},"RuntimeBootstrap":{"version":1,"core_url":"https://core.example/api/v1","device_id":"55555555-5555-4555-8555-555555555555","credential":"fixture-only"},"Deadline":"2099-01-01T00:00:00Z"},"valid":true}, {"kind":"request","name":"reference-null","payload":{"Bootstrap":{"AllocationID":"33333333-3333-4333-8333-333333333333","AllowedDomains":null,"CoreURL":"https://core.example/api/v1","Credential":"fixture-only","DeviceID":"55555555-5555-4555-8555-555555555555","EnvironmentID":"22222222-2222-4222-8222-222222222222","NetworkAccess":"enabled","SessionID":"44444444-4444-4444-8444-444444444444","TenantID":"11111111-1111-4111-8111-111111111111"},"Config":{"APIKey":"","APIURL":"","Binary":"","Domain":"","InstallationID":"66666666-6666-4666-8666-666666666666","Resources":null,"StateDir":"","Template":"","TimeoutSeconds":0},"Deadline":"2099-01-01T00:00:00Z","Operation":"create","Reference":null,"RuntimeBootstrap":{"core_url":"https://core.example/api/v1","credential":"fixture-only","device_id":"55555555-5555-4555-8555-555555555555","version":1},"Version":1},"valid":false}, {"kind":"request","name":"reference-type","payload":{"Bootstrap":{"AllocationID":"33333333-3333-4333-8333-333333333333","AllowedDomains":null,"CoreURL":"https://core.example/api/v1","Credential":"fixture-only","DeviceID":"55555555-5555-4555-8555-555555555555","EnvironmentID":"22222222-2222-4222-8222-222222222222","NetworkAccess":"enabled","SessionID":"44444444-4444-4444-8444-444444444444","TenantID":"11111111-1111-4111-8111-111111111111"},"Config":{"APIKey":"","APIURL":"","Binary":"","Domain":"","InstallationID":"66666666-6666-4666-8666-666666666666","Resources":null,"StateDir":"","Template":"","TimeoutSeconds":0},"Deadline":"2099-01-01T00:00:00Z","Operation":"create","Reference":"invalid","RuntimeBootstrap":{"core_url":"https://core.example/api/v1","credential":"fixture-only","device_id":"55555555-5555-4555-8555-555555555555","version":1},"Version":1},"valid":false}, {"kind":"request","name":"renew","payload":{"Version":1,"Operation":"renew","Config":{"Binary":"","StateDir":"","InstallationID":"66666666-6666-4666-8666-666666666666","APIKey":"","Template":"","APIURL":"","Domain":"","TimeoutSeconds":0,"Resources":null},"Reference":{"TenantID":"11111111-1111-4111-8111-111111111111","EnvironmentID":"22222222-2222-4222-8222-222222222222","AllocationID":"33333333-3333-4333-8333-333333333333"},"Bootstrap":{"TenantID":"11111111-1111-4111-8111-111111111111","EnvironmentID":"22222222-2222-4222-8222-222222222222","AllocationID":"33333333-3333-4333-8333-333333333333","SessionID":"44444444-4444-4444-8444-444444444444","DeviceID":"55555555-5555-4555-8555-555555555555","CoreURL":"https://core.example/api/v1","Credential":"fixture-only","NetworkAccess":"enabled","AllowedDomains":null},"RuntimeBootstrap":{"version":1,"core_url":"https://core.example/api/v1","device_id":"55555555-5555-4555-8555-555555555555","credential":"fixture-only"},"Deadline":"2099-01-01T00:00:00Z"},"valid":true}, +{"kind":"request","name":"resume","payload":{"Version":1,"Operation":"resume","Config":{"Binary":"","StateDir":"","InstallationID":"66666666-6666-4666-8666-666666666666","APIKey":"","Template":"","APIURL":"","Domain":"","TimeoutSeconds":0,"Resources":null},"Reference":{"TenantID":"11111111-1111-4111-8111-111111111111","EnvironmentID":"22222222-2222-4222-8222-222222222222","AllocationID":"33333333-3333-4333-8333-333333333333"},"Bootstrap":{"TenantID":"11111111-1111-4111-8111-111111111111","EnvironmentID":"22222222-2222-4222-8222-222222222222","AllocationID":"33333333-3333-4333-8333-333333333333","SessionID":"44444444-4444-4444-8444-444444444444","DeviceID":"55555555-5555-4555-8555-555555555555","CoreURL":"https://core.example/api/v1","Credential":"fixture-only","NetworkAccess":"enabled","AllowedDomains":null},"RuntimeBootstrap":{"version":1,"core_url":"https://core.example/api/v1","device_id":"55555555-5555-4555-8555-555555555555","credential":"fixture-only"},"Deadline":"2099-01-01T00:00:00Z"},"valid":true}, {"kind":"request","name":"validate_deployment","payload":{"Version":1,"Operation":"validate_deployment","Config":{"Binary":"","StateDir":"","InstallationID":"66666666-6666-4666-8666-666666666666","APIKey":"","Template":"","APIURL":"","Domain":"","TimeoutSeconds":0,"Resources":null},"Reference":{"TenantID":"11111111-1111-4111-8111-111111111111","EnvironmentID":"22222222-2222-4222-8222-222222222222","AllocationID":"33333333-3333-4333-8333-333333333333"},"Bootstrap":{"TenantID":"11111111-1111-4111-8111-111111111111","EnvironmentID":"22222222-2222-4222-8222-222222222222","AllocationID":"33333333-3333-4333-8333-333333333333","SessionID":"44444444-4444-4444-8444-444444444444","DeviceID":"55555555-5555-4555-8555-555555555555","CoreURL":"https://core.example/api/v1","Credential":"fixture-only","NetworkAccess":"enabled","AllowedDomains":null},"RuntimeBootstrap":{"version":1,"core_url":"https://core.example/api/v1","device_id":"55555555-5555-4555-8555-555555555555","credential":"fixture-only"},"Deadline":"2099-01-01T00:00:00Z"},"valid":true}, {"kind":"request","name":"verify_credential","payload":{"Version":1,"Operation":"verify_credential","Config":{"Binary":"","StateDir":"","InstallationID":"66666666-6666-4666-8666-666666666666","APIKey":"","Template":"","APIURL":"","Domain":"","TimeoutSeconds":0,"Resources":null},"Reference":{"TenantID":"11111111-1111-4111-8111-111111111111","EnvironmentID":"22222222-2222-4222-8222-222222222222","AllocationID":"33333333-3333-4333-8333-333333333333"},"References":[{"TenantID":"11111111-1111-4111-8111-111111111111","EnvironmentID":"22222222-2222-4222-8222-222222222222","AllocationID":"33333333-3333-4333-8333-333333333333"}],"Bootstrap":{"TenantID":"11111111-1111-4111-8111-111111111111","EnvironmentID":"22222222-2222-4222-8222-222222222222","AllocationID":"33333333-3333-4333-8333-333333333333","SessionID":"44444444-4444-4444-8444-444444444444","DeviceID":"55555555-5555-4555-8555-555555555555","CoreURL":"https://core.example/api/v1","Credential":"fixture-only","NetworkAccess":"enabled","AllowedDomains":null},"RuntimeBootstrap":{"version":1,"core_url":"https://core.example/api/v1","device_id":"55555555-5555-4555-8555-555555555555","credential":"fixture-only"},"Deadline":"2099-01-01T00:00:00Z"},"valid":true}, {"kind":"request","name":"verify_credential-count-0","payload":{"Version":1,"Operation":"verify_credential","Config":{"Binary":"","StateDir":"","InstallationID":"66666666-6666-4666-8666-666666666666","APIKey":"","Template":"","APIURL":"","Domain":"","TimeoutSeconds":0,"Resources":null},"Reference":{"TenantID":"11111111-1111-4111-8111-111111111111","EnvironmentID":"22222222-2222-4222-8222-222222222222","AllocationID":"33333333-3333-4333-8333-333333333333"},"Bootstrap":{"TenantID":"11111111-1111-4111-8111-111111111111","EnvironmentID":"22222222-2222-4222-8222-222222222222","AllocationID":"33333333-3333-4333-8333-333333333333","SessionID":"44444444-4444-4444-8444-444444444444","DeviceID":"55555555-5555-4555-8555-555555555555","CoreURL":"https://core.example/api/v1","Credential":"fixture-only","NetworkAccess":"enabled","AllowedDomains":null},"RuntimeBootstrap":{"version":1,"core_url":"https://core.example/api/v1","device_id":"55555555-5555-4555-8555-555555555555","credential":"fixture-only"},"Deadline":"2099-01-01T00:00:00Z"},"valid":true}, diff --git a/services/core/tools/microsandbox-provider/bootstrap.go b/services/core/tools/microsandbox-provider/bootstrap.go index ad1af890a..dfceb93c4 100644 --- a/services/core/tools/microsandbox-provider/bootstrap.go +++ b/services/core/tools/microsandbox-provider/bootstrap.go @@ -5,6 +5,7 @@ package main import ( "context" "encoding/json" + "github.com/MiniMax-AI/OpenAgentCore/internal/runtimebootstrap" "time" "github.com/MiniMax-AI/OpenAgentCore/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