From 78355d491ec3add54a3cbfc4d6697937a8c562f2 Mon Sep 17 00:00:00 2001 From: saladday <1203511142@qq.com> Date: Wed, 30 Sep 2026 05:43:26 +0000 Subject: [PATCH 1/2] refactor: require explicit provider operation contracts --- CONTRIBUTING.md | 3 + apps/docs/content/docs/sandbox-provider.mdx | 65 ++++++-- apps/docs/content/guide-sources.json | 4 +- .../agents-api/node-generation-protocol.md | 16 ++ docs/sandbox-provider.md | 65 ++++++-- .../server/managed_generation_operations.go | 149 ++++++++++++++++++ .../cmd/server/managed_generations.go | 35 ++-- .../agents-api/cmd/server/managed_setup.go | 32 ++-- .../archive_cancellation_cleanup_test.go | 10 +- .../provider_operations_fixture_test.go | 70 ++++++++ .../internal/execution/runtime_compute.go | 12 +- .../internal/execution/runtime_lifecycle.go | 5 +- .../internal/providercontract/operations.go | 95 +++++++++++ .../runtimeobs/operations_fixture_test.go | 28 ++++ .../internal/runtimeobs/operations_test.go | 74 +++++++++ .../agents-api/internal/runtimeobs/service.go | 48 +++++- .../internal/runtimeobs/service_test.go | 4 +- .../agents-api/internal/runtimeobs/source.go | 13 +- .../internal/sandbox/contracttest/provider.go | 8 + .../internal/sandbox/docker/operations.go | 74 +++++++++ .../internal/sandbox/e2b/observations.go | 6 +- .../internal/sandbox/e2b/observations_test.go | 6 +- .../internal/sandbox/e2b/operations.go | 65 ++++++++ .../sandbox/microsandbox/operations.go | 47 ++++++ .../agents-api/internal/sandbox/node/agent.go | 11 +- .../internal/sandbox/node/generations.go | 4 +- .../agents-api/internal/sandbox/node/hub.go | 2 +- .../internal/sandbox/node/node_test.go | 2 +- .../internal/sandbox/node/observations.go | 6 +- .../sandbox/node/observations_test.go | 7 +- .../internal/sandbox/node/operations.go | 30 ++++ .../internal/sandbox/node/operations_test.go | 58 +++++++ .../node/provider_operations_fixture_test.go | 92 +++++++++++ .../agents-api/internal/sandbox/node/proxy.go | 60 +++++-- .../agents-api/internal/sandbox/node/wire.go | 56 +++++-- .../agents-api/internal/sandbox/operations.go | 80 ++++++++++ .../internal/sandbox/operations_test.go | 112 +++++++++++++ .../internal/sandbox/providers/config.go | 4 + .../internal/sandbox/providers/e2b.go | 9 +- .../internal/sandbox/providers/operations.go | 18 +++ .../sandbox/providers/operations_test.go | 22 +++ .../internal/sandbox/providers/registry.go | 27 ++-- .../sandbox/providers/registry_test.go | 10 +- .../internal/sandbox/sandbox_provider.go | 13 +- .../store/provider_operations_fixture_test.go | 92 +++++++++++ 45 files changed, 1508 insertions(+), 141 deletions(-) create mode 100644 services/agents-api/cmd/server/managed_generation_operations.go create mode 100644 services/agents-api/internal/execution/provider_operations_fixture_test.go create mode 100644 services/agents-api/internal/providercontract/operations.go create mode 100644 services/agents-api/internal/runtimeobs/operations_fixture_test.go create mode 100644 services/agents-api/internal/runtimeobs/operations_test.go create mode 100644 services/agents-api/internal/sandbox/docker/operations.go create mode 100644 services/agents-api/internal/sandbox/e2b/operations.go create mode 100644 services/agents-api/internal/sandbox/microsandbox/operations.go create mode 100644 services/agents-api/internal/sandbox/node/operations.go create mode 100644 services/agents-api/internal/sandbox/node/operations_test.go create mode 100644 services/agents-api/internal/sandbox/node/provider_operations_fixture_test.go create mode 100644 services/agents-api/internal/sandbox/operations.go create mode 100644 services/agents-api/internal/sandbox/operations_test.go create mode 100644 services/agents-api/internal/sandbox/providers/operations.go create mode 100644 services/agents-api/internal/sandbox/providers/operations_test.go create mode 100644 services/agents-api/internal/store/provider_operations_fixture_test.go diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index ce8974c85..f46401340 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -203,6 +203,9 @@ The Runtime is not a sandbox; see implementation must not require a new orchestration path selected by its name. Sandbox registration, configuration adaptation and persistence boundaries follow the [Sandbox Provider guide](docs/sandbox-provider.md#register-the-provider-kind). + Resource operation declarations are exhaustive and validated against the existing + small interfaces; support is never inferred from method presence. See the + [explicit operation contract](docs/sandbox-provider.md#explicit-operation-contracts). - Core owns durable Session/Turn state and scheduling. Runtime owns local execution resources. Harness adapters translate the common execution contract into native operations; model and sandbox provider details stay behind their diff --git a/apps/docs/content/docs/sandbox-provider.mdx b/apps/docs/content/docs/sandbox-provider.mdx index 8e221c313..c36582f5e 100644 --- a/apps/docs/content/docs/sandbox-provider.mdx +++ b/apps/docs/content/docs/sandbox-provider.mdx @@ -35,8 +35,8 @@ connectivity; see [Runtime and outer isolation](/concepts#runtime-and-outer-isol ## Steps -1. **Read the contract.** Implement the five required operations and any - optional interfaces in [Implement the interface](#implement-the-interface). +1. **Read the contract.** Implement the five required operations and explicitly handle all + extension interfaces in [Implement the interface](#implement-the-interface). 2. **Write the adapter package** under `services/agents-api/internal/sandbox/` (native SDK calls, ownership checks, identity translation, private config). Assert `var _ sandbox.SandboxProvider = (*YourAdapter)(nil)` at compile time. @@ -69,20 +69,53 @@ Use the existing types; do not introduce another lifecycle protocol or a vendor-specific execution path. A backend without a native renewable lease (Docker) still keeps service-owned hosted expiry and cleanup requirements. -### Required and optional interfaces +### Explicit operation contracts + +Keep the existing small interfaces. Every Provider must implement their methods +and return a complete `ProviderOperations()` declaration. The interface methods +are the operation inventory; `sandbox.ValidateOperations` checks it without a +second hand-maintained list. | Contract | Requirement | Responsibility | | --- | --- | --- | -| `sandbox.SandboxProvider` | Required | `Create`, `GetInfo`, `Renew`, `Kill`, and bounded bootstrap/diagnostic `RunCommand` | -| `sandbox.CheckpointProvider` | Optional, separate interface | Exact compute incarnations, snapshot capture/restore, retained-source resume and cleanup | -| `runtimeobs.Source` | Optional, separate interface | Read-only, ownership-checked resource observations | -| `runtimeobs.BatchSource` | Optional, separate interface | Bounded observations in input order, with per-target errors; `ok=false` means no batch read occurred | - -Do not implement an optional interface with successful no-op methods. Core uses -interface assertions to select optional operations. Advertised support requires -contract and native acceptance evidence; a healthy node or an available CLI does -not establish it. Observation never renews a lease, starts compute or prepares a -Harness. See the [observation contract](https://github.com/MiniMax-AI/parsar-core/blob/f6d258735fc601c521dd990e6f9e1ed261f4ef2d/contracts/agents-api/runtime-observability.md). +| `sandbox.SandboxProvider` | Five operations must be supported | Allocation lifecycle and bounded administrative commands | +| `sandbox.CheckpointProvider` | Explicit supported or unsupported decision for every method | Exact compute incarnations, capture/restore, retained-source resume and cleanup | +| `runtimeobs.Source` | Explicit decision | Ownership-checked read-only observations | +| `runtimeobs.BatchSource` | Explicit decision | 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 | + +Each declaration entry has `state: supported` with no reason, or +`state: unsupported` with an authored reason code. Missing entries, zero states, +unknown entries, missing methods and unsafe reasons fail validation. The entire +checkpoint lifecycle must agree on support; batch observation requires single +observation. Adding a method to an existing interface requires an explicit +decision and implementation in every adapter. Never supply a default base class +or generate blanket unsupported implementations for future methods. + +Unsupported methods return `providercontract.UnsupportedError` before native +I/O. The error identifies the exact operation and a safe code, not a native +message, resource identity, endpoint or credential. Empty results, nil errors, +`Unavailable`, and unknown mutation outcomes cannot substitute for unsupported. +The five required methods cannot return unsupported; a backend without a native +lease preserves the existing read-only `Renew` semantics. + +Each adapter owns one `Operations()` function, shared by its concrete instance +and registration. `providers.ValidateBinding` checks both against the existing +interfaces and against each other. Runtime admission and node generation loading +also reject incomplete providers. Interface assertions establish method shape +only; callers use the declaration to decide whether an operation is supported. + +`ObserveBatch` returns an error instead of an ambiguous boolean. Only a typed, +safe `UnsupportedError` for `ObserveBatch` permits per-target `Observe` calls. +An unavailable service, timeout or other failure never triggers that fallback. +Observation never renews, starts, prepares or stops compute; see the +[observation contract](https://github.com/MiniMax-AI/parsar-core/blob/f6d258735fc601c521dd990e6f9e1ed261f4ef2d/contracts/agents-api/runtime-observability.md). + +Contract tests call every declared unsupported native method with no configured +native client, require its matching error and zero result, and reject incomplete +or contradictory declarations. Supported behavior still requires native and +lifecycle tests; declaration validation alone cannot prove SDK semantics. `RunCommand` is an existing administrative bootstrap/diagnostic facility, not an alternate route for Skills, Plugins, MCP setup, initial files, Session execution or @@ -173,20 +206,20 @@ A provider must not implement a competing preparation path. `sandbox/providers/registry.go` is the sole registration table. Each entry binds an adapter's specification/resource validators, selection normalization, deployment -mode, defaults, optional checkpoint capability, and local or direct constructor. +mode, defaults, the adapter-owned operation declaration, and local or direct constructor. `providers.Build` constructs node-local adapters; `providers.BuildDirect` constructs direct adapters. Neither allocates compute. There is no init-time registration or runtime plugin loading. For a new implementation: -1. Implement `SandboxProvider` in its adapter package and add native contract tests. +1. Implement the operation contracts above in the adapter package and add native contract tests. 2. Add its configuration validators and optional read-only `SelectionDiscoverer` for native resource discovery. Normalization must copy input before changing it. `RestoreSelection` must retain access to owned resources without requiring new template validation. Put native credential verification behind `CredentialVerifier` when needed. -3. Register its constructor, policies and defaults in `providers/registry.go`. +3. Register its constructor, policies, operation declaration and defaults in `providers/registry.go`. Node proxy identity and checkpoint advertisement consume this same entry. The installer projection uses those registered policies and the common field bounds in `sandbox/deployment_contract.go`; regenerate it with diff --git a/apps/docs/content/guide-sources.json b/apps/docs/content/guide-sources.json index 6c9c5fc66..eb5a3b903 100644 --- a/apps/docs/content/guide-sources.json +++ b/apps/docs/content/guide-sources.json @@ -28,7 +28,7 @@ "contracts/agents-api/harness-onboarding.md": "874a80dc1ffc5e97b2783bea3f6b217893f997b38bb8180a29f09c6dbf75e622", "docs/runtime-bootstrap.md": "0d49aed73b298039e04453e6f465b0e925fb2227fa35d4206820bb3acbd8df39", "docs/runtime-protocol.md": "e8aa4cf862b5f63a4138cfeb25196a6bdb5c40a2e389e828b464585415dc6c45", - "docs/sandbox-provider.md": "4d47b234a457f874f3ff7e61a7f5da6fc8b75b0c6df0ce0ce9ead1c07afdfb3e", + "docs/sandbox-provider.md": "a1323c7f6a28a4f5f0e823274c08422707a5dab85b11e008ccb28a3aafc02058", "apps/docs/scripts/guides.json": "3c3768fb94fd3d42464d4ce8c55724b9ba8360133c7cc19d513d8ec83a8e034f" }, "outputs": { @@ -59,6 +59,6 @@ "content/docs/harness-onboarding.mdx": "17626e9f6256ddfffc95fd81dfa8d9c712d9ff16792ea5ef54bfe8c3ce1a2a9f", "content/docs/runtime-bootstrap.mdx": "58982955811c9ad46a9fc04d8d2aef5a762fc9d25d5833f88aa75ee3fc46523f", "content/docs/runtime-protocol.mdx": "d6b818edf2a0e0c4f2d3878f75805043c6dba3b09c6563f80faa4cecd05d40c7", - "content/docs/sandbox-provider.mdx": "9db833fe9ff3f30a65451f80d7251348695a08e7830f9e95871739a710529818" + "content/docs/sandbox-provider.mdx": "83a6244fb2140836aa6746c990c9cd04d04e58d5b3a52faf1a84e5e846dfd99f" } } diff --git a/contracts/agents-api/node-generation-protocol.md b/contracts/agents-api/node-generation-protocol.md index 00f5dbba9..122463aac 100644 --- a/contracts/agents-api/node-generation-protocol.md +++ b/contracts/agents-api/node-generation-protocol.md @@ -10,6 +10,22 @@ Every Provider request carries its exact allocation-owned deployment generation; Core never strips it for an older peer. Automatic preparation and retention frames are sent only to nodes advertising generation management. +## Provider operation outcomes + +Node startup and generation loading validate complete Provider operation +declarations before accepting work. Proxies use the same registered declaration +for admission; unsupported operations reject before node resolution or native I/O. +The Provider [operation contract](../../docs/sandbox-provider.md#explicit-operation-contracts) +owns the inventory and support rules. + +A Provider response with `error_code: unsupported` carries an `unsupported` object +containing the exact method `operation` and an authored safe `reason` code. The +proxy checks both against the request. Missing, malformed or mismatched evidence +is an unconfirmed result, never proof that a mutation was rejected. Unsupported +remains distinct from observation unavailability and unknown compute/command +results; it does not settle resource ownership or authorize replay. The current +private wire version requires both peers to understand this outcome. + ## Bounded control A generation-managing node's hello or heartbeat contains at most eight generation observations. diff --git a/docs/sandbox-provider.md b/docs/sandbox-provider.md index 60a6af2d7..c59a14371 100644 --- a/docs/sandbox-provider.md +++ b/docs/sandbox-provider.md @@ -32,8 +32,8 @@ connectivity; see [Runtime and outer isolation](design-principles.md#runtime-and ## Steps -1. **Read the contract.** Implement the five required operations and any - optional interfaces in [Implement the interface](#implement-the-interface). +1. **Read the contract.** Implement the five required operations and explicitly handle all + extension interfaces in [Implement the interface](#implement-the-interface). 2. **Write the adapter package** under `services/agents-api/internal/sandbox/` (native SDK calls, ownership checks, identity translation, private config). Assert `var _ sandbox.SandboxProvider = (*YourAdapter)(nil)` at compile time. @@ -66,20 +66,53 @@ Use the existing types; do not introduce another lifecycle protocol or a vendor-specific execution path. A backend without a native renewable lease (Docker) still keeps service-owned hosted expiry and cleanup requirements. -### Required and optional interfaces +### Explicit operation contracts + +Keep the existing small interfaces. Every Provider must implement their methods +and return a complete `ProviderOperations()` declaration. The interface methods +are the operation inventory; `sandbox.ValidateOperations` checks it without a +second hand-maintained list. | Contract | Requirement | Responsibility | | --- | --- | --- | -| `sandbox.SandboxProvider` | Required | `Create`, `GetInfo`, `Renew`, `Kill`, and bounded bootstrap/diagnostic `RunCommand` | -| `sandbox.CheckpointProvider` | Optional, separate interface | Exact compute incarnations, snapshot capture/restore, retained-source resume and cleanup | -| `runtimeobs.Source` | Optional, separate interface | Read-only, ownership-checked resource observations | -| `runtimeobs.BatchSource` | Optional, separate interface | Bounded observations in input order, with per-target errors; `ok=false` means no batch read occurred | - -Do not implement an optional interface with successful no-op methods. Core uses -interface assertions to select optional operations. Advertised support requires -contract and native acceptance evidence; a healthy node or an available CLI does -not establish it. Observation never renews a lease, starts compute or prepares a -Harness. See the [observation contract](../contracts/agents-api/runtime-observability.md). +| `sandbox.SandboxProvider` | Five operations must be supported | Allocation lifecycle and bounded administrative commands | +| `sandbox.CheckpointProvider` | Explicit supported or unsupported decision for every method | Exact compute incarnations, capture/restore, retained-source resume and cleanup | +| `runtimeobs.Source` | Explicit decision | Ownership-checked read-only observations | +| `runtimeobs.BatchSource` | Explicit decision | 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 | + +Each declaration entry has `state: supported` with no reason, or +`state: unsupported` with an authored reason code. Missing entries, zero states, +unknown entries, missing methods and unsafe reasons fail validation. The entire +checkpoint lifecycle must agree on support; batch observation requires single +observation. Adding a method to an existing interface requires an explicit +decision and implementation in every adapter. Never supply a default base class +or generate blanket unsupported implementations for future methods. + +Unsupported methods return `providercontract.UnsupportedError` before native +I/O. The error identifies the exact operation and a safe code, not a native +message, resource identity, endpoint or credential. Empty results, nil errors, +`Unavailable`, and unknown mutation outcomes cannot substitute for unsupported. +The five required methods cannot return unsupported; a backend without a native +lease preserves the existing read-only `Renew` semantics. + +Each adapter owns one `Operations()` function, shared by its concrete instance +and registration. `providers.ValidateBinding` checks both against the existing +interfaces and against each other. Runtime admission and node generation loading +also reject incomplete providers. Interface assertions establish method shape +only; callers use the declaration to decide whether an operation is supported. + +`ObserveBatch` returns an error instead of an ambiguous boolean. Only a typed, +safe `UnsupportedError` for `ObserveBatch` permits per-target `Observe` calls. +An unavailable service, timeout or other failure never triggers that fallback. +Observation never renews, starts, prepares or stops compute; see the +[observation contract](../contracts/agents-api/runtime-observability.md). + +Contract tests call every declared unsupported native method with no configured +native client, require its matching error and zero result, and reject incomplete +or contradictory declarations. Supported behavior still requires native and +lifecycle tests; declaration validation alone cannot prove SDK semantics. `RunCommand` is an existing administrative bootstrap/diagnostic facility, not an alternate route for Skills, Plugins, MCP setup, initial files, Session execution or @@ -170,20 +203,20 @@ A provider must not implement a competing preparation path. `sandbox/providers/registry.go` is the sole registration table. Each entry binds an adapter's specification/resource validators, selection normalization, deployment -mode, defaults, optional checkpoint capability, and local or direct constructor. +mode, defaults, the adapter-owned operation declaration, and local or direct constructor. `providers.Build` constructs node-local adapters; `providers.BuildDirect` constructs direct adapters. Neither allocates compute. There is no init-time registration or runtime plugin loading. For a new implementation: -1. Implement `SandboxProvider` in its adapter package and add native contract tests. +1. Implement the operation contracts above in the adapter package and add native contract tests. 2. Add its configuration validators and optional read-only `SelectionDiscoverer` for native resource discovery. Normalization must copy input before changing it. `RestoreSelection` must retain access to owned resources without requiring new template validation. Put native credential verification behind `CredentialVerifier` when needed. -3. Register its constructor, policies and defaults in `providers/registry.go`. +3. Register its constructor, policies, operation declaration and defaults in `providers/registry.go`. Node proxy identity and checkpoint advertisement consume this same entry. The installer projection uses those registered policies and the common field bounds in `sandbox/deployment_contract.go`; regenerate it with diff --git a/services/agents-api/cmd/server/managed_generation_operations.go b/services/agents-api/cmd/server/managed_generation_operations.go new file mode 100644 index 000000000..3416cfa85 --- /dev/null +++ b/services/agents-api/cmd/server/managed_generation_operations.go @@ -0,0 +1,149 @@ +package main + +import ( + "context" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" +) + +func (p *generationRouter) Initial(ctx context.Context, r sandbox.Reference) (sandbox.Compute, error) { + if err := providercontract.Require(p, "Initial"); err != nil { + return sandbox.Compute{}, err + } + v, done, err := p.route(ctx, r) + if err != nil { + return sandbox.Compute{}, err + } + defer done() + cp, err := sandbox.Checkpoint(v) + if err != nil { + return sandbox.Compute{}, err + } + return cp.Initial(ctx, r) +} +func (p *generationRouter) NewCompute(ctx context.Context, r sandbox.Reference, g uint64, snapshot *sandbox.SnapshotIdentity) (sandbox.Compute, error) { + if err := providercontract.Require(p, "NewCompute"); err != nil { + return sandbox.Compute{}, err + } + v, done, err := p.route(ctx, r) + if err != nil { + return sandbox.Compute{}, err + } + defer done() + cp, err := sandbox.Checkpoint(v) + if err != nil { + return sandbox.Compute{}, err + } + return cp.NewCompute(ctx, r, g, snapshot) +} +func (p *generationRouter) GetCompute(ctx context.Context, r sandbox.Reference, c sandbox.Compute) (sandbox.ComputeState, error) { + if err := providercontract.Require(p, "GetCompute"); err != nil { + return sandbox.ComputeState{}, err + } + v, done, err := p.route(ctx, r) + if err != nil { + return sandbox.ComputeState{}, err + } + defer done() + cp, err := sandbox.Checkpoint(v) + if err != nil { + return sandbox.ComputeState{}, err + } + return cp.GetCompute(ctx, r, c) +} +func (p *generationRouter) Suspend(ctx context.Context, q sandbox.SuspendRequest) (sandbox.ComputeState, error) { + if err := providercontract.Require(p, "Suspend"); err != nil { + return sandbox.ComputeState{}, err + } + v, done, err := p.route(ctx, q.Reference) + if err != nil { + return sandbox.ComputeState{}, err + } + defer done() + cp, err := sandbox.Checkpoint(v) + if err != nil { + return sandbox.ComputeState{}, err + } + return cp.Suspend(ctx, q) +} +func (p *generationRouter) Resume(ctx context.Context, q sandbox.ResumeRequest) (sandbox.ComputeState, error) { + if err := providercontract.Require(p, "Resume"); err != nil { + return sandbox.ComputeState{}, err + } + v, done, err := p.route(ctx, q.Reference) + if err != nil { + return sandbox.ComputeState{}, err + } + defer done() + cp, err := sandbox.Checkpoint(v) + if err != nil { + return sandbox.ComputeState{}, err + } + return cp.Resume(ctx, q) +} +func (p *generationRouter) KillCompute(ctx context.Context, r sandbox.Reference, c sandbox.Compute) error { + if err := providercontract.Require(p, "KillCompute"); err != nil { + return err + } + v, done, err := p.route(ctx, r) + if err != nil { + return err + } + defer done() + cp, err := sandbox.Checkpoint(v) + if err != nil { + return err + } + return cp.KillCompute(ctx, r, c) +} +func (p *generationRouter) DeleteSnapshot(ctx context.Context, r sandbox.Reference, snapshot sandbox.SnapshotIdentity) error { + if err := providercontract.Require(p, "DeleteSnapshot"); err != nil { + return err + } + v, done, err := p.route(ctx, r) + if err != nil { + return err + } + defer done() + cp, err := sandbox.Checkpoint(v) + if err != nil { + return err + } + return cp.DeleteSnapshot(ctx, r, snapshot) +} +func (p *generationRouter) RunCommandCompute(ctx context.Context, r sandbox.Reference, c sandbox.Compute, command sandbox.Command) (sandbox.CommandResult, error) { + if err := providercontract.Require(p, "RunCommandCompute"); err != nil { + return sandbox.CommandResult{}, err + } + v, done, err := p.route(ctx, r) + if err != nil { + return sandbox.CommandResult{}, err + } + defer done() + cp, err := sandbox.Checkpoint(v) + if err != nil { + return sandbox.CommandResult{}, err + } + return cp.RunCommandCompute(ctx, r, c, command) +} +func (p *generationRouter) ResumeCompute(ctx context.Context, r sandbox.Reference, c sandbox.Compute) (sandbox.ComputeState, error) { + if err := providercontract.Require(p, "ResumeCompute"); err != nil { + return sandbox.ComputeState{}, err + } + v, done, err := p.route(ctx, r) + if err != nil { + return sandbox.ComputeState{}, err + } + defer done() + cp, err := sandbox.Checkpoint(v) + if err != nil { + return sandbox.ComputeState{}, err + } + return cp.ResumeCompute(ctx, r, c) +} +func (*generationRouter) DiscoverSelection(context.Context, sandbox.Selection) (sandbox.Selection, error) { + return sandbox.Selection{}, &providercontract.UnsupportedError{Operation: "DiscoverSelection", Reason: "generation_router_does_not_discover_configuration"} +} +func (*generationRouter) VerifyCredential(context.Context, []sandbox.Reference) error { + return &providercontract.UnsupportedError{Operation: "VerifyCredential", Reason: "generation_router_does_not_verify_configuration"} +} diff --git a/services/agents-api/cmd/server/managed_generations.go b/services/agents-api/cmd/server/managed_generations.go index cb3ba82c2..23a265f06 100644 --- a/services/agents-api/cmd/server/managed_generations.go +++ b/services/agents-api/cmd/server/managed_generations.go @@ -3,6 +3,8 @@ package main import ( "context" "errors" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "maps" "time" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/execution" @@ -23,8 +25,9 @@ type generationStore interface { // immutable specification and current credential. There is no mutable provider // map to unload and no current-generation fallback for a missing historical row. type generationRouter struct { - setup *managedSetup - store generationStore + setup *managedSetup + store generationStore + operations providercontract.Operations } func (p *generationRouter) route(ctx context.Context, r sandbox.Reference) (sandbox.SandboxProvider, func(), error) { @@ -87,15 +90,22 @@ func (p *generationRouter) RunCommand(ctx context.Context, r sandbox.Reference, type observedGenerationRouter struct{ *generationRouter } +func (p *generationRouter) ProviderOperations() providercontract.Operations { + return maps.Clone(p.operations) +} + func (p *observedGenerationRouter) Observe(ctx context.Context, t runtimeobs.Target) (runtimeobs.Sample, error) { v, done, err := p.route(ctx, sandbox.Reference{TenantID: t.TenantID, EnvironmentID: t.EnvironmentID, AllocationID: t.Instance.AllocationID}) if err != nil { return runtimeobs.Sample{}, err } defer done() + if err := providercontract.Require(v, "Observe"); err != nil { + return runtimeobs.Sample{}, err + } source, ok := v.(runtimeobs.Source) if !ok { - return runtimeobs.Sample{}, runtimeobs.ErrUnavailable + return runtimeobs.Sample{}, providercontract.ErrContract } return source.Observe(ctx, t) } @@ -112,11 +122,13 @@ func (s *managedSetup) routeGenerations(candidate execution.PreparedRuntimeDeplo if !ok { return execution.PreparedRuntimeDeployment{}, errors.New("sandbox generation store is unavailable") } - router := &generationRouter{setup: s, store: db} - if _, ok := candidate.Config.Provider.(runtimeobs.Source); ok { - candidate.Config.Provider = &observedGenerationRouter{router} - } else { - candidate.Config.Provider = router + router := &generationRouter{setup: s, store: db, operations: candidate.Config.Provider.ProviderOperations()} + router.operations["ObserveBatch"] = providercontract.Support{State: providercontract.Unsupported, Reason: "allocations_require_individual_generation_routing"} + router.operations["DiscoverSelection"] = providercontract.Support{State: providercontract.Unsupported, Reason: "generation_router_does_not_discover_configuration"} + router.operations["VerifyCredential"] = providercontract.Support{State: providercontract.Unsupported, Reason: "generation_router_does_not_verify_configuration"} + candidate.Config.Provider = &observedGenerationRouter{router} + if err := sandbox.ValidateProvider(candidate.Config.Provider); err != nil { + return execution.PreparedRuntimeDeployment{}, err } if !adapter.Credential { return candidate, nil @@ -139,6 +151,9 @@ func (s *managedSetup) routeGenerations(candidate execution.PreparedRuntimeDeplo if err != nil { return err } + if err := providercontract.Require(provider, "VerifyCredential"); err != nil { + return err + } verifier, ok := provider.(sandbox.CredentialVerifier) if !ok { return sandbox.ErrInvalid @@ -221,6 +236,6 @@ func (s *managedSetup) routeGenerations(candidate execution.PreparedRuntimeDeplo // A batch can contain allocations from different endpoint generations. Let the // observation worker route each allocation through its immutable generation. -func (p *observedGenerationRouter) ObserveBatch(_ context.Context, _ []runtimeobs.Target) ([]runtimeobs.BatchResult, bool) { - return nil, false +func (p *observedGenerationRouter) ObserveBatch(_ context.Context, _ []runtimeobs.Target) ([]runtimeobs.BatchResult, error) { + return nil, &providercontract.UnsupportedError{Operation: "ObserveBatch", Reason: "allocations_require_individual_generation_routing"} } diff --git a/services/agents-api/cmd/server/managed_setup.go b/services/agents-api/cmd/server/managed_setup.go index 2ab0c49f6..2b632f911 100644 --- a/services/agents-api/cmd/server/managed_setup.go +++ b/services/agents-api/cmd/server/managed_setup.go @@ -4,6 +4,7 @@ import ( "context" "errors" "fmt" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" "sync/atomic" "time" @@ -94,7 +95,8 @@ func (s *managedSetup) prepare(ctx context.Context, setup store.SandboxSetup) (e return execution.PreparedRuntimeDeployment{}, err } selection := sandbox.Selection{Provider: setup.Provider, DeploymentSpec: setup.Specification, E2B: setup.E2B} - if discoverer, ok := candidate.Config.Provider.(sandbox.SelectionDiscoverer); ok { + if err := providercontract.Require(candidate.Config.Provider, "DiscoverSelection"); err == nil { + discoverer := candidate.Config.Provider.(sandbox.SelectionDiscoverer) selection, err = discoverer.DiscoverSelection(ctx, selection) if err != nil { return execution.PreparedRuntimeDeployment{}, err @@ -121,13 +123,16 @@ func (s *managedSetup) configuration(setup store.SandboxSetup) (execution.Prepar } selected := &execution.RuntimeProvider{InstallationID: setup.InstallationID, ProviderKind: setup.Provider, Generation: setup.Generation, Mode: setup.Mode, AdmissionPaused: setup.AdmissionPaused, CoreURL: s.publicURL + "/api/v1", BackendFingerprint: setup.BackendFingerprint, Provider: provider} - if _, supportsCheckpoint := provider.(sandbox.CheckpointProvider); supportsCheckpoint { + if sandbox.SupportsCheckpoint(provider) { selected.Suspension = &execution.RuntimeSuspensionPolicy{IdleTimeout: time.Duration(setup.IdleSeconds) * time.Second, Retention: time.Duration(setup.RetentionSeconds) * time.Second, MaxActive: 4, MaxRetained: 16} } return execution.PreparedRuntimeDeployment{Config: selected, Publish: s.publish}, nil } +func (*managedSetup) ProviderOperations() providercontract.Operations { + return providercontract.Operations{"Observe": {State: providercontract.Supported}, "ObserveBatch": {State: providercontract.Supported}} +} func (s *managedSetup) ObservationProviderType() string { if selected := s.selected.Load(); selected != nil && selected.Config != nil { return selected.Config.ProviderKind @@ -135,16 +140,22 @@ func (s *managedSetup) ObservationProviderType() string { return "" } -// ObserveBatch delegates to the selected provider's batch read. It reports -// ok=false for providers without one, so each target uses Observe instead. -func (s *managedSetup) ObserveBatch(ctx context.Context, targets []runtimeobs.Target) ([]runtimeobs.BatchResult, bool) { +// ObserveBatch preserves the selected provider's explicit Unsupported or failure. +// Only Unsupported permits the observation service to read targets individually. +func (s *managedSetup) ObserveBatch(ctx context.Context, targets []runtimeobs.Target) ([]runtimeobs.BatchResult, error) { selected, err := s.load(ctx) - if err != nil || selected == nil { - return nil, false + if err != nil { + return nil, err + } + if selected == nil { + return nil, runtimeobs.ErrUnavailable + } + if err := providercontract.Require(selected.Provider, "ObserveBatch"); err != nil { + return nil, err } source, ok := selected.Provider.(runtimeobs.BatchSource) if !ok { - return nil, false + return nil, providercontract.ErrContract } return source.ObserveBatch(ctx, targets) } @@ -157,9 +168,12 @@ func (s *managedSetup) Observe(ctx context.Context, target runtimeobs.Target) (r if selected == nil { return runtimeobs.Sample{}, runtimeobs.ErrUnavailable } + if err := providercontract.Require(selected.Provider, "Observe"); err != nil { + return runtimeobs.Sample{}, err + } source, ok := selected.Provider.(runtimeobs.Source) if !ok { - return runtimeobs.Sample{}, runtimeobs.ErrUnavailable + return runtimeobs.Sample{}, providercontract.ErrContract } return source.Observe(ctx, target) } diff --git a/services/agents-api/internal/execution/archive_cancellation_cleanup_test.go b/services/agents-api/internal/execution/archive_cancellation_cleanup_test.go index 9470d7f89..7b4ebfde1 100644 --- a/services/agents-api/internal/execution/archive_cancellation_cleanup_test.go +++ b/services/agents-api/internal/execution/archive_cancellation_cleanup_test.go @@ -3,23 +3,25 @@ package execution import ( "context" "encoding/json" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox/microsandbox" "net/http" "net/http/httptest" "net/url" "testing" "time" + v1 "github.com/MiniMax-AI-Dev/parsar/contracts/agents-api/v1" "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/device" "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/gateway" "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/proto" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/adminaudit" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/identity" + runtimegateway "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtime" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" "github.com/google/uuid" "github.com/gorilla/websocket" - runtimegateway "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtime" - v1 "github.com/MiniMax-AI-Dev/parsar/contracts/agents-api/v1" ) // Provider callbacks inspect the real database at the instant destructive @@ -37,6 +39,10 @@ func (p waitingCleanupProvider) Kill(context.Context, sandbox.Reference) error { return nil } +func (waitingCleanupCheckpoint) ProviderOperations() providercontract.Operations { + return microsandbox.Operations() +} + type waitingCleanupCheckpoint struct { sandbox.CheckpointProvider beforeKill func() diff --git a/services/agents-api/internal/execution/provider_operations_fixture_test.go b/services/agents-api/internal/execution/provider_operations_fixture_test.go new file mode 100644 index 000000000..624068305 --- /dev/null +++ b/services/agents-api/internal/execution/provider_operations_fixture_test.go @@ -0,0 +1,70 @@ +package execution + +import ( + "context" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeobs" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" +) + +func (*lifecycleOnlySandbox) ProviderOperations() providercontract.Operations { + return providercontract.Operations{ + "Create": {State: providercontract.Supported}, + "GetInfo": {State: providercontract.Supported}, + "Renew": {State: providercontract.Supported}, + "Kill": {State: providercontract.Supported}, + "RunCommand": {State: providercontract.Supported}, + "Initial": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "NewCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "GetCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "Suspend": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "Resume": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "KillCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "DeleteSnapshot": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "RunCommandCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "ResumeCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "Observe": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "ObserveBatch": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "DiscoverSelection": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "VerifyCredential": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + } +} +func (*lifecycleOnlySandbox) Initial(context.Context, sandbox.Reference) (sandbox.Compute, error) { + return sandbox.Compute{}, &providercontract.UnsupportedError{Operation: "Initial", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleOnlySandbox) NewCompute(context.Context, sandbox.Reference, uint64, *sandbox.SnapshotIdentity) (sandbox.Compute, error) { + return sandbox.Compute{}, &providercontract.UnsupportedError{Operation: "NewCompute", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleOnlySandbox) GetCompute(context.Context, sandbox.Reference, sandbox.Compute) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "GetCompute", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleOnlySandbox) Suspend(context.Context, sandbox.SuspendRequest) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "Suspend", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleOnlySandbox) Resume(context.Context, sandbox.ResumeRequest) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "Resume", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleOnlySandbox) KillCompute(context.Context, sandbox.Reference, sandbox.Compute) error { + return &providercontract.UnsupportedError{Operation: "KillCompute", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleOnlySandbox) DeleteSnapshot(context.Context, sandbox.Reference, sandbox.SnapshotIdentity) error { + return &providercontract.UnsupportedError{Operation: "DeleteSnapshot", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleOnlySandbox) RunCommandCompute(context.Context, sandbox.Reference, sandbox.Compute, sandbox.Command) (sandbox.CommandResult, error) { + return sandbox.CommandResult{}, &providercontract.UnsupportedError{Operation: "RunCommandCompute", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleOnlySandbox) ResumeCompute(context.Context, sandbox.Reference, sandbox.Compute) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "ResumeCompute", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleOnlySandbox) Observe(context.Context, runtimeobs.Target) (runtimeobs.Sample, error) { + return runtimeobs.Sample{}, &providercontract.UnsupportedError{Operation: "Observe", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleOnlySandbox) ObserveBatch(context.Context, []runtimeobs.Target) ([]runtimeobs.BatchResult, error) { + return nil, &providercontract.UnsupportedError{Operation: "ObserveBatch", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleOnlySandbox) DiscoverSelection(context.Context, sandbox.Selection) (sandbox.Selection, error) { + return sandbox.Selection{}, &providercontract.UnsupportedError{Operation: "DiscoverSelection", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleOnlySandbox) VerifyCredential(context.Context, []sandbox.Reference) error { + return &providercontract.UnsupportedError{Operation: "VerifyCredential", Reason: "fixture_operation_not_supported"} +} diff --git a/services/agents-api/internal/execution/runtime_compute.go b/services/agents-api/internal/execution/runtime_compute.go index 7fc56384a..eb4ed7804 100644 --- a/services/agents-api/internal/execution/runtime_compute.go +++ b/services/agents-api/internal/execution/runtime_compute.go @@ -63,9 +63,9 @@ func (r *runtimeLifecycle) saveCompute(ctx context.Context, owner store.RuntimeA return r.store.SetRuntimeCompute(ctx, owner, phase, raw, until, idleTimeout) } func (r *runtimeLifecycle) enableCompute(ctx context.Context, owner store.RuntimeAllocation) error { - p, ok := r.config.Provider.(sandbox.CheckpointProvider) - if !ok { - return sandbox.ErrInvalid + p, capabilityErr := sandbox.Checkpoint(r.config.Provider) + if capabilityErr != nil { + return capabilityErr } initial, err := p.Initial(ctx, runtimeReference(owner)) if err != nil { @@ -83,9 +83,9 @@ func (r *runtimeLifecycle) enableCompute(ctx context.Context, owner store.Runtim } func (r *runtimeLifecycle) observeCompute(ctx context.Context, owner store.RuntimeAllocation) error { - p, ok := r.config.Provider.(sandbox.CheckpointProvider) - if !ok { - return sandbox.ErrInvalid + p, capabilityErr := sandbox.Checkpoint(r.config.Provider) + if capabilityErr != nil { + return capabilityErr } var state runtimeCompute if json.Unmarshal(owner.ComputeState, &state) != nil || state.Current.ID == "" { diff --git a/services/agents-api/internal/execution/runtime_lifecycle.go b/services/agents-api/internal/execution/runtime_lifecycle.go index d372a4df9..b084a7daa 100644 --- a/services/agents-api/internal/execution/runtime_lifecycle.go +++ b/services/agents-api/internal/execution/runtime_lifecycle.go @@ -86,6 +86,9 @@ func validatedRuntimeProvider(config *RuntimeProvider, registry *gateway.Registr if err != nil || id == uuid.Nil || id.String() != config.InstallationID || fingerprintErr != nil || len(fingerprint) != 32 || hex.EncodeToString(fingerprint) != config.BackendFingerprint { return RuntimeProvider{}, sandbox.ErrInvalid } + if err := sandbox.ValidateProvider(config.Provider); err != nil { + return RuntimeProvider{}, err + } copied := *config if copied.Mode == "" && copied.ProviderKind != "" { copied.Mode = "nodes" @@ -98,7 +101,7 @@ func validatedRuntimeProvider(config *RuntimeProvider, registry *gateway.Registr } if config.Suspension != nil { policy := *config.Suspension - if _, ok := config.Provider.(sandbox.CheckpointProvider); !ok || policy.IdleTimeout < time.Second || policy.Retention < time.Second || policy.MaxActive < 1 || policy.MaxRetained < policy.MaxActive { + if !sandbox.SupportsCheckpoint(config.Provider) || policy.IdleTimeout < time.Second || policy.Retention < time.Second || policy.MaxActive < 1 || policy.MaxRetained < policy.MaxActive { return RuntimeProvider{}, sandbox.ErrInvalid } copied.Suspension = &policy diff --git a/services/agents-api/internal/providercontract/operations.go b/services/agents-api/internal/providercontract/operations.go new file mode 100644 index 000000000..ab1aa7a32 --- /dev/null +++ b/services/agents-api/internal/providercontract/operations.go @@ -0,0 +1,95 @@ +// Package providercontract describes explicit resource-provider operations. +// It does not describe Harness capabilities or external API optional fields. +package providercontract + +import ( + "errors" + "fmt" + "reflect" + "regexp" +) + +var ErrUnsupported = errors.New("provider operation unsupported") +var ErrContract = errors.New("invalid provider operation declaration") +var reasonPattern = regexp.MustCompile(`^[a-z][a-z0-9_]{0,95}$`) + +type Support struct { + State string `json:"state"` + Reason string `json:"reason,omitempty"` +} + +const Supported = "supported" +const Unsupported = "unsupported" + +// Zero values are omissions, never an unsupported declaration. +type Operations map[string]Support +type Declared interface{ ProviderOperations() Operations } + +// UnsupportedError contains only an operation name and an authored safe code. +// It never carries a native error, credential, endpoint or resource identifier. +type UnsupportedError struct { + Operation string `json:"operation"` + Reason string `json:"reason"` +} + +func (e *UnsupportedError) Error() string { + return "provider operation " + e.Operation + " unsupported: " + e.Reason +} +func (e *UnsupportedError) Unwrap() error { return ErrUnsupported } + +func (s Support) Check(operation string) error { + switch { + case s.State == Supported && s.Reason == "": + return nil + case s.State == Unsupported && reasonPattern.MatchString(s.Reason): + return &UnsupportedError{Operation: operation, Reason: s.Reason} + default: + return fmt.Errorf("%w: %s", ErrContract, operation) + } +} + +// UnsupportedReason verifies the operation and safe code before callers fall back +// or publish the reason. A bare sentinel is a contract violation. +func UnsupportedReason(err error, operation string) (string, bool) { + var value *UnsupportedError + if !errors.As(err, &value) || value.Operation != operation || !reasonPattern.MatchString(value.Reason) { + return "", false + } + return value.Reason, true +} + +func Require(p Declared, operation string) error { + if p == nil { + return ErrContract + } + return p.ProviderOperations()[operation].Check(operation) +} + +// Validate checks the existing small interfaces, including methods declared +// unsupported. Every adapter must explicitly implement those rejections. +// A new method cannot be silently covered by an old declaration or stub. +func Validate(p Declared, contracts ...reflect.Type) error { + if p == nil { + return ErrContract + } + value := reflect.ValueOf(p) + if value.Kind() == reflect.Pointer && value.IsNil() { + return ErrContract + } + operations := p.ProviderOperations() + for _, contract := range contracts { + if !value.Type().Implements(contract) { + return fmt.Errorf("%w: missing %s implementation", ErrContract, contract.Name()) + } + for i := 0; i < contract.NumMethod(); i++ { + name := contract.Method(i).Name + if name == "ProviderOperations" { + continue + } + if err := operations[name].Check(name); err != nil && !errors.Is(err, ErrUnsupported) { + return err + } + } + } + return nil +} diff --git a/services/agents-api/internal/runtimeobs/operations_fixture_test.go b/services/agents-api/internal/runtimeobs/operations_fixture_test.go new file mode 100644 index 000000000..1ac5df001 --- /dev/null +++ b/services/agents-api/internal/runtimeobs/operations_fixture_test.go @@ -0,0 +1,28 @@ +package runtimeobs + +import ( + "context" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" +) + +func (*fixedSource) ProviderOperations() providercontract.Operations { + return providercontract.Operations{"Observe": {State: providercontract.Supported}, "ObserveBatch": {State: providercontract.Unsupported, Reason: "fixture_has_no_batch_observation"}} +} +func (*fixedSource) ObserveBatch(context.Context, []Target) ([]BatchResult, error) { + return nil, &providercontract.UnsupportedError{Operation: "ObserveBatch", Reason: "fixture_has_no_batch_observation"} +} +func (blockingSource) ProviderOperations() providercontract.Operations { + return providercontract.Operations{"Observe": {State: providercontract.Supported}, "ObserveBatch": {State: providercontract.Unsupported, Reason: "fixture_has_no_batch_observation"}} +} +func (blockingSource) ObserveBatch(context.Context, []Target) ([]BatchResult, error) { + return nil, &providercontract.UnsupportedError{Operation: "ObserveBatch", Reason: "fixture_has_no_batch_observation"} +} +func (*countingSource) ProviderOperations() providercontract.Operations { + return providercontract.Operations{"Observe": {State: providercontract.Supported}, "ObserveBatch": {State: providercontract.Unsupported, Reason: "fixture_has_no_batch_observation"}} +} +func (*countingSource) ObserveBatch(context.Context, []Target) ([]BatchResult, error) { + return nil, &providercontract.UnsupportedError{Operation: "ObserveBatch", Reason: "fixture_has_no_batch_observation"} +} +func (*batchSource) ProviderOperations() providercontract.Operations { + return providercontract.Operations{"Observe": {State: providercontract.Supported}, "ObserveBatch": {State: providercontract.Supported}} +} diff --git a/services/agents-api/internal/runtimeobs/operations_test.go b/services/agents-api/internal/runtimeobs/operations_test.go new file mode 100644 index 000000000..50f30ea98 --- /dev/null +++ b/services/agents-api/internal/runtimeobs/operations_test.go @@ -0,0 +1,74 @@ +package runtimeobs + +import ( + "context" + "errors" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "testing" + "time" +) + +type failingBatchSource struct { + *fixedSource + batchErr error + batches int +} + +func (*failingBatchSource) ProviderOperations() providercontract.Operations { + return providercontract.Operations{"Observe": {State: providercontract.Supported}, "ObserveBatch": {State: providercontract.Supported}} +} +func (s *failingBatchSource) ObserveBatch(context.Context, []Target) ([]BatchResult, error) { + s.batches++ + return nil, s.batchErr +} +func TestBatchFallbackRequiresExplicitSafeUnsupported(t *testing.T) { + for _, test := range []struct { + name string + err error + fallback bool + }{ + {"unsupported", &providercontract.UnsupportedError{Operation: "ObserveBatch", Reason: "native_batch_not_supported"}, true}, + {"unavailable", ErrUnavailable, false}, + {"timeout", context.DeadlineExceeded, false}, + {"failure", errors.New("provider failed"), false}, + {"bare unsupported", providercontract.ErrUnsupported, false}, + {"unsafe unsupported", &providercontract.UnsupportedError{Operation: "ObserveBatch", Reason: "private endpoint / key"}, false}, + {"wrong operation", &providercontract.UnsupportedError{Operation: "Observe", Reason: "not_supported"}, false}, + } { + t.Run(test.name, func(t *testing.T) { + now := time.Now() + ratio := 0.5 + source := &failingBatchSource{fixedSource: &fixedSource{sample: Sample{ObservedAt: now, CPUUtilizationRatio: &ratio}}, batchErr: test.err} + target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} + service, err := NewService(fixedResolver{target: target}, map[string]Source{"provider": source}) + if err != nil { + t.Fatal(err) + } + observations, errs := service.ObserveSessions(t.Context(), []SessionIdentity{{TenantID: "tenant", SessionID: "session"}}, PageOptions{}) + if source.batches != 1 || (source.calls == 1) != test.fallback { + t.Fatalf("batch=%d single=%d", source.batches, source.calls) + } + if test.fallback && (errs[0] != nil || observations[0].Status != StatusObserved) { + t.Fatal(observations, errs) + } + }) + } +} + +type unsupportedObservation struct{ *fixedSource } + +func (*unsupportedObservation) ProviderOperations() providercontract.Operations { + return providercontract.Operations{"Observe": {State: providercontract.Unsupported, Reason: "native_metrics_not_supported"}, "ObserveBatch": {State: providercontract.Unsupported, Reason: "native_metrics_not_supported"}} +} +func TestUnsupportedObservationIsNotUnavailable(t *testing.T) { + source := &unsupportedObservation{&fixedSource{}} + target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} + service, err := NewService(fixedResolver{target: target}, map[string]Source{"provider": source}) + if err != nil { + t.Fatal(err) + } + observation, err := service.ObserveSession(t.Context(), "tenant", "session") + if err != nil || observation.Status != StatusUnsupported || observation.Reason != "native_metrics_not_supported" || source.calls != 0 { + t.Fatal(observation, err, source.calls) + } +} diff --git a/services/agents-api/internal/runtimeobs/service.go b/services/agents-api/internal/runtimeobs/service.go index 79b28dc31..08188e606 100644 --- a/services/agents-api/internal/runtimeobs/service.go +++ b/services/agents-api/internal/runtimeobs/service.go @@ -4,6 +4,8 @@ import ( "context" "errors" "fmt" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "reflect" "regexp" "sync" "time" @@ -53,6 +55,9 @@ func NewService(resolver TargetResolver, sources map[string]Source, options ...S if key == "" || source == nil { return nil, errors.New("invalid Runtime observation source") } + if err := providercontract.Validate(source, reflect.TypeFor[Source](), reflect.TypeFor[BatchSource]()); err != nil { + return nil, err + } copySources[key] = source } service := &Service{resolver: resolver, sources: copySources, now: time.Now} @@ -108,7 +113,7 @@ type PageOptions struct { } // ObserveSessions observes one page of Sessions. Running targets of a source -// that implements BatchSource share one provider read per MaxBatchTargets; +// that explicitly supports BatchSource share one provider read per MaxBatchTargets; // other sources are read per target, exactly as ObserveSession reads them. // Results and errors are aligned with sessions. func (s *Service) ObserveSessions(ctx context.Context, sessions []SessionIdentity, options PageOptions) ([]Observation, []error) { @@ -179,7 +184,11 @@ func (s *Service) observeSessions(ctx context.Context, sessions []SessionIdentit read := chunk[index] sourceCtx, stop := sourceContext(ctx, options.SourceTimeout) started := time.Now() - sample, err := read.source.Observe(sourceCtx, read.target) + var sample Sample + err := providercontract.Require(read.source, "Observe") + if err == nil { + sample, err = read.source.Observe(sourceCtx, read.target) + } stop() observations[read.index], errs[read.index] = s.complete(ctx, read, sample, err, time.Since(started), collectionSource, owner) }) @@ -192,8 +201,20 @@ func (s *Service) observeSessions(ctx context.Context, sessions []SessionIdentit // for its current provider. func (s *Service) readBatch(ctx context.Context, chunk []*sourceRead, observations []Observation, errs []error, collectionSource CollectionSource, owner OwnershipChecker, sourceTimeout time.Duration) bool { batch, ok := chunk[0].source.(BatchSource) - if !ok { - return false + if !ok { // NewService rejects this; never treat malformed registration as unsupported. + for _, read := range chunk { + errs[read.index] = providercontract.ErrContract + } + return true + } + if err := providercontract.Require(chunk[0].source, "ObserveBatch"); err != nil { + if errors.Is(err, providercontract.ErrUnsupported) { + return false + } + for _, read := range chunk { + errs[read.index] = err + } + return true } targets := make([]Target, len(chunk)) for index, read := range chunk { @@ -206,11 +227,20 @@ func (s *Service) readBatch(ctx context.Context, chunk []*sourceRead, observatio } sourceCtx, stop := sourceContext(ctx, sourceTimeout) started := time.Now() - results, ok := batch.ObserveBatch(sourceCtx, targets) + results, batchErr := batch.ObserveBatch(sourceCtx, targets) stop() - if !ok { + if _, unsupported := providercontract.UnsupportedReason(batchErr, "ObserveBatch"); unsupported { return false } + if errors.Is(batchErr, providercontract.ErrUnsupported) { + batchErr = providercontract.ErrContract + } + if batchErr != nil { + for _, read := range chunk { + observations[read.index], errs[read.index] = s.complete(ctx, read, Sample{}, batchErr, 0, collectionSource, owner) + } + return true + } duration := time.Since(started) for index, read := range chunk { if len(results) != len(chunk) { @@ -288,6 +318,12 @@ func (s *Service) resolve(ctx context.Context, tenantID, sessionID string, owner func (s *Service) complete(ctx context.Context, read *sourceRead, sample Sample, err error, sourceDuration time.Duration, collectionSource CollectionSource, owner OwnershipChecker) (Observation, error) { observation := Observation{Target: read.target, Status: StatusUnavailable, ProviderType: read.providerType, SourceDuration: sourceDuration} switch { + case errors.Is(err, providercontract.ErrUnsupported): + reason, valid := providercontract.UnsupportedReason(err, "Observe") + if !valid { + return Observation{}, providercontract.ErrContract + } + observation.Status, observation.Reason = StatusUnsupported, reason case errors.Is(err, context.DeadlineExceeded): observation.Reason = "sample_timeout" case errors.Is(err, ErrNotRunning): diff --git a/services/agents-api/internal/runtimeobs/service_test.go b/services/agents-api/internal/runtimeobs/service_test.go index 30db2529f..51d5b9e59 100644 --- a/services/agents-api/internal/runtimeobs/service_test.go +++ b/services/agents-api/internal/runtimeobs/service_test.go @@ -624,7 +624,7 @@ func (s *batchSource) Observe(context.Context, Target) (Sample, error) { return Sample{}, errors.New("per-target read used for a batch source") } -func (s *batchSource) ObserveBatch(ctx context.Context, targets []Target) ([]BatchResult, bool) { +func (s *batchSource) ObserveBatch(ctx context.Context, targets []Target) ([]BatchResult, error) { time.Sleep(time.Millisecond) s.mu.Lock() s.batches = append(s.batches, len(targets)) @@ -636,7 +636,7 @@ func (s *batchSource) ObserveBatch(ctx context.Context, targets []Target) ([]Bat results[index].Err = ErrNotRunning } } - return results, true + return results, nil } func TestServiceBatchesPageReadsWithinProviderLimit(t *testing.T) { diff --git a/services/agents-api/internal/runtimeobs/source.go b/services/agents-api/internal/runtimeobs/source.go index f52e8aa2b..fcfb9b1f7 100644 --- a/services/agents-api/internal/runtimeobs/source.go +++ b/services/agents-api/internal/runtimeobs/source.go @@ -3,6 +3,7 @@ package runtimeobs import ( "context" "errors" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" ) var ( @@ -13,19 +14,19 @@ var ( // Source reads one provider-owned Runtime instance. Implementations must verify // ownership before returning data and must not renew, restart, or stop compute. type Source interface { + providercontract.Declared Observe(context.Context, Target) (Sample, error) } // MaxBatchTargets bounds the targets of one BatchSource read. const MaxBatchTargets = 100 -// BatchSource is an optional Source extension for providers that read many -// Runtime instances in one bounded request. ObserveBatch returns ok=false, -// without reading, when the current provider has no batch read; the Service -// then reads each target with Observe. Otherwise it returns one result per -// target, in order, under the same ownership rules as Observe. +// BatchSource reads multiple instances under the same ownership rules as Observe. +// Unsupported returns a typed providercontract.UnsupportedError before reading; +// only that result permits per-target observation. Unavailable or failed reads +// must not silently retry through Observe. Successful results preserve target order. type BatchSource interface { - ObserveBatch(context.Context, []Target) (results []BatchResult, ok bool) + ObserveBatch(context.Context, []Target) (results []BatchResult, err error) } type BatchResult struct { diff --git a/services/agents-api/internal/sandbox/contracttest/provider.go b/services/agents-api/internal/sandbox/contracttest/provider.go index f72a8ee79..8c7555f53 100644 --- a/services/agents-api/internal/sandbox/contracttest/provider.go +++ b/services/agents-api/internal/sandbox/contracttest/provider.go @@ -4,6 +4,8 @@ package contracttest import ( "context" + "errors" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" "reflect" "testing" @@ -51,6 +53,9 @@ func RunFailures(t *testing.T, factory Factory) { ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second) defer cancel() f := factory(t, Scenario{operation, fault}, cancel) + if err := sandbox.ValidateProvider(f.Provider); err != nil { + t.Fatal(err) + } var info sandbox.Info var err error switch operation { @@ -63,6 +68,9 @@ func RunFailures(t *testing.T, factory Factory) { case "kill": err = f.Provider.Kill(ctx, f.Bootstrap.Reference) } + if errors.Is(err, providercontract.ErrUnsupported) { + t.Fatal("required operation returned Unsupported", err) + } if err == nil { t.Fatal("native failure was reported as success") } diff --git a/services/agents-api/internal/sandbox/docker/operations.go b/services/agents-api/internal/sandbox/docker/operations.go new file mode 100644 index 000000000..1a6b9dd2c --- /dev/null +++ b/services/agents-api/internal/sandbox/docker/operations.go @@ -0,0 +1,74 @@ +package docker + +import ( + "context" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeobs" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" +) + +var _ sandbox.CheckpointProvider = (*Provider)(nil) +var _ sandbox.SelectionDiscoverer = (*Provider)(nil) +var _ sandbox.CredentialVerifier = (*Provider)(nil) +var _ runtimeobs.BatchSource = (*Provider)(nil) + +// Operations is this adapter's complete authored resource contract. +func Operations() providercontract.Operations { + return providercontract.Operations{ + "Create": {State: providercontract.Supported}, + "GetInfo": {State: providercontract.Supported}, + "Renew": {State: providercontract.Supported}, + "Kill": {State: providercontract.Supported}, + "RunCommand": {State: providercontract.Supported}, + "Initial": {State: providercontract.Unsupported, Reason: "docker_does_not_support_checkpoints"}, + "NewCompute": {State: providercontract.Unsupported, Reason: "docker_does_not_support_checkpoints"}, + "GetCompute": {State: providercontract.Unsupported, Reason: "docker_does_not_support_checkpoints"}, + "Suspend": {State: providercontract.Unsupported, Reason: "docker_does_not_support_checkpoints"}, + "Resume": {State: providercontract.Unsupported, Reason: "docker_does_not_support_checkpoints"}, + "KillCompute": {State: providercontract.Unsupported, Reason: "docker_does_not_support_checkpoints"}, + "DeleteSnapshot": {State: providercontract.Unsupported, Reason: "docker_does_not_support_checkpoints"}, + "RunCommandCompute": {State: providercontract.Unsupported, Reason: "docker_does_not_support_checkpoints"}, + "ResumeCompute": {State: providercontract.Unsupported, Reason: "docker_does_not_support_checkpoints"}, + "Observe": {State: providercontract.Supported}, + "ObserveBatch": {State: providercontract.Unsupported, Reason: "docker_does_not_support_batch_observation"}, + "DiscoverSelection": {State: providercontract.Unsupported, Reason: "docker_does_not_support_selection_discovery"}, + "VerifyCredential": {State: providercontract.Unsupported, Reason: "docker_does_not_support_credential_verification"}, + } +} +func (*Provider) ProviderOperations() providercontract.Operations { return Operations() } +func (p *Provider) Initial(context.Context, sandbox.Reference) (sandbox.Compute, error) { + return sandbox.Compute{}, &providercontract.UnsupportedError{Operation: "Initial", Reason: Operations()["Initial"].Reason} +} +func (p *Provider) NewCompute(context.Context, sandbox.Reference, uint64, *sandbox.SnapshotIdentity) (sandbox.Compute, error) { + return sandbox.Compute{}, &providercontract.UnsupportedError{Operation: "NewCompute", Reason: Operations()["NewCompute"].Reason} +} +func (p *Provider) GetCompute(context.Context, sandbox.Reference, sandbox.Compute) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "GetCompute", Reason: Operations()["GetCompute"].Reason} +} +func (p *Provider) Suspend(context.Context, sandbox.SuspendRequest) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "Suspend", Reason: Operations()["Suspend"].Reason} +} +func (p *Provider) Resume(context.Context, sandbox.ResumeRequest) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "Resume", Reason: Operations()["Resume"].Reason} +} +func (p *Provider) KillCompute(context.Context, sandbox.Reference, sandbox.Compute) error { + return &providercontract.UnsupportedError{Operation: "KillCompute", Reason: Operations()["KillCompute"].Reason} +} +func (p *Provider) DeleteSnapshot(context.Context, sandbox.Reference, sandbox.SnapshotIdentity) error { + return &providercontract.UnsupportedError{Operation: "DeleteSnapshot", Reason: Operations()["DeleteSnapshot"].Reason} +} +func (p *Provider) RunCommandCompute(context.Context, sandbox.Reference, sandbox.Compute, sandbox.Command) (sandbox.CommandResult, error) { + return sandbox.CommandResult{}, &providercontract.UnsupportedError{Operation: "RunCommandCompute", Reason: Operations()["RunCommandCompute"].Reason} +} +func (p *Provider) ResumeCompute(context.Context, sandbox.Reference, sandbox.Compute) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "ResumeCompute", Reason: Operations()["ResumeCompute"].Reason} +} +func (p *Provider) ObserveBatch(context.Context, []runtimeobs.Target) ([]runtimeobs.BatchResult, error) { + return nil, &providercontract.UnsupportedError{Operation: "ObserveBatch", Reason: Operations()["ObserveBatch"].Reason} +} +func (p *Provider) DiscoverSelection(context.Context, sandbox.Selection) (sandbox.Selection, error) { + return sandbox.Selection{}, &providercontract.UnsupportedError{Operation: "DiscoverSelection", Reason: Operations()["DiscoverSelection"].Reason} +} +func (p *Provider) VerifyCredential(context.Context, []sandbox.Reference) error { + return &providercontract.UnsupportedError{Operation: "VerifyCredential", Reason: Operations()["VerifyCredential"].Reason} +} diff --git a/services/agents-api/internal/sandbox/e2b/observations.go b/services/agents-api/internal/sandbox/e2b/observations.go index 6d7edef7c..e083bef97 100644 --- a/services/agents-api/internal/sandbox/e2b/observations.go +++ b/services/agents-api/internal/sandbox/e2b/observations.go @@ -43,7 +43,7 @@ func (p *Provider) Observe(ctx context.Context, target runtimeobs.Target) (runti // helper request. The helper takes sandbox IDs from its private receipts, // confirms each running sandbox by its allocation labels and reads E2B's batch // metrics once. It never connects to, renews or changes a sandbox. -func (p *Provider) ObserveBatch(ctx context.Context, targets []runtimeobs.Target) ([]runtimeobs.BatchResult, bool) { +func (p *Provider) ObserveBatch(ctx context.Context, targets []runtimeobs.Target) ([]runtimeobs.BatchResult, error) { results := make([]runtimeobs.BatchResult, len(targets)) references := make([]sandbox.Reference, 0, len(targets)) positions := make([]int, 0, len(targets)) @@ -60,7 +60,7 @@ func (p *Provider) ObserveBatch(ctx context.Context, targets []runtimeobs.Target } } if len(references) == 0 { - return results, true + return results, nil } observations, err := p.observe(ctx, references) now := p.now() @@ -71,7 +71,7 @@ func (p *Provider) ObserveBatch(ctx context.Context, targets []runtimeobs.Target } results[index].Sample, results[index].Err = sampleFromObservation(observations[offset], now) } - return results, true + return results, nil } func (p *Provider) observe(ctx context.Context, references []sandbox.Reference) ([]Observation, error) { diff --git a/services/agents-api/internal/sandbox/e2b/observations_test.go b/services/agents-api/internal/sandbox/e2b/observations_test.go index 54f96c4da..771e3dd23 100644 --- a/services/agents-api/internal/sandbox/e2b/observations_test.go +++ b/services/agents-api/internal/sandbox/e2b/observations_test.go @@ -27,10 +27,10 @@ func TestObserveBatchMapsMetricsAndKeepsUnmeasuredValuesNull(t *testing.T) { targets = append(targets, runtimeobs.Target{TenantID: reference.TenantID, EnvironmentID: reference.EnvironmentID, Mode: runtimeobs.ModeManaged, Instance: runtimeobs.Instance{AllocationID: reference.AllocationID, ProviderKey: p.config.InstallationID}}) } - results, ok := p.ObserveBatch(bounded(t), targets) - if !ok || len(caller.requests) != 1 || caller.requests[0].Operation != "observe" || len(caller.requests[0].References) != 2 || + results, err := p.ObserveBatch(bounded(t), targets) + if err != nil || len(caller.requests) != 1 || caller.requests[0].Operation != "observe" || len(caller.requests[0].References) != 2 || caller.requests[0].References[1] != stopped || caller.requests[0].Deadline.IsZero() { - t.Fatalf("batch was not one bounded helper request: ok=%v %+v", ok, caller.requests) + t.Fatalf("batch was not one bounded helper request: err=%v %+v", err, caller.requests) } sample := results[0].Sample if results[0].Err != nil || !sample.ObservedAt.Equal(now) || !sample.StartedAt.Equal(started) || diff --git a/services/agents-api/internal/sandbox/e2b/operations.go b/services/agents-api/internal/sandbox/e2b/operations.go new file mode 100644 index 000000000..ee303792f --- /dev/null +++ b/services/agents-api/internal/sandbox/e2b/operations.go @@ -0,0 +1,65 @@ +package e2b + +import ( + "context" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeobs" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" +) + +var _ sandbox.CheckpointProvider = (*Provider)(nil) +var _ sandbox.SelectionDiscoverer = (*Provider)(nil) +var _ sandbox.CredentialVerifier = (*Provider)(nil) +var _ runtimeobs.BatchSource = (*Provider)(nil) + +// Operations is this adapter's complete authored resource contract. +func Operations() providercontract.Operations { + return providercontract.Operations{ + "Create": {State: providercontract.Supported}, + "GetInfo": {State: providercontract.Supported}, + "Renew": {State: providercontract.Supported}, + "Kill": {State: providercontract.Supported}, + "RunCommand": {State: providercontract.Supported}, + "Initial": {State: providercontract.Unsupported, Reason: "e2b_does_not_support_checkpoints"}, + "NewCompute": {State: providercontract.Unsupported, Reason: "e2b_does_not_support_checkpoints"}, + "GetCompute": {State: providercontract.Unsupported, Reason: "e2b_does_not_support_checkpoints"}, + "Suspend": {State: providercontract.Unsupported, Reason: "e2b_does_not_support_checkpoints"}, + "Resume": {State: providercontract.Unsupported, Reason: "e2b_does_not_support_checkpoints"}, + "KillCompute": {State: providercontract.Unsupported, Reason: "e2b_does_not_support_checkpoints"}, + "DeleteSnapshot": {State: providercontract.Unsupported, Reason: "e2b_does_not_support_checkpoints"}, + "RunCommandCompute": {State: providercontract.Unsupported, Reason: "e2b_does_not_support_checkpoints"}, + "ResumeCompute": {State: providercontract.Unsupported, Reason: "e2b_does_not_support_checkpoints"}, + "Observe": {State: providercontract.Supported}, + "ObserveBatch": {State: providercontract.Supported}, + "DiscoverSelection": {State: providercontract.Supported}, + "VerifyCredential": {State: providercontract.Supported}, + } +} +func (*Provider) ProviderOperations() providercontract.Operations { return Operations() } +func (p *Provider) Initial(context.Context, sandbox.Reference) (sandbox.Compute, error) { + return sandbox.Compute{}, &providercontract.UnsupportedError{Operation: "Initial", Reason: Operations()["Initial"].Reason} +} +func (p *Provider) NewCompute(context.Context, sandbox.Reference, uint64, *sandbox.SnapshotIdentity) (sandbox.Compute, error) { + return sandbox.Compute{}, &providercontract.UnsupportedError{Operation: "NewCompute", Reason: Operations()["NewCompute"].Reason} +} +func (p *Provider) GetCompute(context.Context, sandbox.Reference, sandbox.Compute) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "GetCompute", Reason: Operations()["GetCompute"].Reason} +} +func (p *Provider) Suspend(context.Context, sandbox.SuspendRequest) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "Suspend", Reason: Operations()["Suspend"].Reason} +} +func (p *Provider) Resume(context.Context, sandbox.ResumeRequest) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "Resume", Reason: Operations()["Resume"].Reason} +} +func (p *Provider) KillCompute(context.Context, sandbox.Reference, sandbox.Compute) error { + return &providercontract.UnsupportedError{Operation: "KillCompute", Reason: Operations()["KillCompute"].Reason} +} +func (p *Provider) DeleteSnapshot(context.Context, sandbox.Reference, sandbox.SnapshotIdentity) error { + return &providercontract.UnsupportedError{Operation: "DeleteSnapshot", Reason: Operations()["DeleteSnapshot"].Reason} +} +func (p *Provider) RunCommandCompute(context.Context, sandbox.Reference, sandbox.Compute, sandbox.Command) (sandbox.CommandResult, error) { + return sandbox.CommandResult{}, &providercontract.UnsupportedError{Operation: "RunCommandCompute", Reason: Operations()["RunCommandCompute"].Reason} +} +func (p *Provider) ResumeCompute(context.Context, sandbox.Reference, sandbox.Compute) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "ResumeCompute", Reason: Operations()["ResumeCompute"].Reason} +} diff --git a/services/agents-api/internal/sandbox/microsandbox/operations.go b/services/agents-api/internal/sandbox/microsandbox/operations.go new file mode 100644 index 000000000..cd0014448 --- /dev/null +++ b/services/agents-api/internal/sandbox/microsandbox/operations.go @@ -0,0 +1,47 @@ +package microsandbox + +import ( + "context" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeobs" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" +) + +var _ sandbox.CheckpointProvider = (*Provider)(nil) +var _ sandbox.SelectionDiscoverer = (*Provider)(nil) +var _ sandbox.CredentialVerifier = (*Provider)(nil) +var _ runtimeobs.BatchSource = (*Provider)(nil) + +// Operations is this adapter's complete authored resource contract. +func Operations() providercontract.Operations { + return providercontract.Operations{ + "Create": {State: providercontract.Supported}, + "GetInfo": {State: providercontract.Supported}, + "Renew": {State: providercontract.Supported}, + "Kill": {State: providercontract.Supported}, + "RunCommand": {State: providercontract.Supported}, + "Initial": {State: providercontract.Supported}, + "NewCompute": {State: providercontract.Supported}, + "GetCompute": {State: providercontract.Supported}, + "Suspend": {State: providercontract.Supported}, + "Resume": {State: providercontract.Supported}, + "KillCompute": {State: providercontract.Supported}, + "DeleteSnapshot": {State: providercontract.Supported}, + "RunCommandCompute": {State: providercontract.Supported}, + "ResumeCompute": {State: providercontract.Supported}, + "Observe": {State: providercontract.Supported}, + "ObserveBatch": {State: providercontract.Unsupported, Reason: "microsandbox_does_not_support_batch_observation"}, + "DiscoverSelection": {State: providercontract.Unsupported, Reason: "microsandbox_does_not_support_selection_discovery"}, + "VerifyCredential": {State: providercontract.Unsupported, Reason: "microsandbox_does_not_support_credential_verification"}, + } +} +func (*Provider) ProviderOperations() providercontract.Operations { return Operations() } +func (p *Provider) ObserveBatch(context.Context, []runtimeobs.Target) ([]runtimeobs.BatchResult, error) { + return nil, &providercontract.UnsupportedError{Operation: "ObserveBatch", Reason: Operations()["ObserveBatch"].Reason} +} +func (p *Provider) DiscoverSelection(context.Context, sandbox.Selection) (sandbox.Selection, error) { + return sandbox.Selection{}, &providercontract.UnsupportedError{Operation: "DiscoverSelection", Reason: Operations()["DiscoverSelection"].Reason} +} +func (p *Provider) VerifyCredential(context.Context, []sandbox.Reference) error { + return &providercontract.UnsupportedError{Operation: "VerifyCredential", Reason: Operations()["VerifyCredential"].Reason} +} diff --git a/services/agents-api/internal/sandbox/node/agent.go b/services/agents-api/internal/sandbox/node/agent.go index 7b95f08ed..4883979c3 100644 --- a/services/agents-api/internal/sandbox/node/agent.go +++ b/services/agents-api/internal/sandbox/node/agent.go @@ -3,6 +3,7 @@ package node import ( "context" "errors" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox/providers" "math/rand/v2" "net/http" "net/url" @@ -79,9 +80,13 @@ func Run(ctx context.Context, config AgentConfig) error { if config.Generations == nil && (config.Provider == nil || config.Probe == nil) { return sandbox.ErrInvalid } - _, checkpoint := config.Provider.(sandbox.CheckpointProvider) - if config.Generations == nil && checkpoint != (config.Identity.Provider == "microsandbox") { - return sandbox.ErrInvalid + if config.Generations == nil { + if err := sandbox.ValidateProvider(config.Provider); err != nil { + return err + } + if sandbox.SupportsCheckpoint(config.Provider) != providers.SupportsCheckpoint(config.Identity.Provider) { + return sandbox.ErrInvalid + } } release, err := lockDirectory(config.StateDirectory) if err != nil { diff --git a/services/agents-api/internal/sandbox/node/generations.go b/services/agents-api/internal/sandbox/node/generations.go index 4fa625845..fc5707875 100644 --- a/services/agents-api/internal/sandbox/node/generations.go +++ b/services/agents-api/internal/sandbox/node/generations.go @@ -60,7 +60,7 @@ func NewGenerationManager(ctx context.Context, options GenerationManagerOptions) owned, cancel := context.WithCancel(ctx) m := &GenerationManager{values: map[uint64]*localGeneration{}, options: options, ctx: owned, cancel: cancel, wake: make(chan struct{}, 1)} for _, v := range options.Initial { - if !validGeneration(v.Generation) || !validSpecificationDigest(v.SpecificationDigest) || v.Provider == nil || v.Probe == nil || m.values[v.Generation] != nil { + if !validGeneration(v.Generation) || !validSpecificationDigest(v.SpecificationDigest) || sandbox.ValidateProvider(v.Provider) != nil || v.Probe == nil || m.values[v.Generation] != nil { cancel() return nil, sandbox.ErrInvalid } @@ -291,7 +291,7 @@ func (m *GenerationManager) prepareLoop() { m.preparingCancel = nil g.preparing = false g.refs-- - if err == nil && (value.Generation != generation || value.SpecificationDigest != digest || value.Provider == nil || value.Probe == nil) { + if err == nil && (value.Generation != generation || value.SpecificationDigest != digest || sandbox.ValidateProvider(value.Provider) != nil || value.Probe == nil) { err = sandbox.ErrOwnership } if err == nil { diff --git a/services/agents-api/internal/sandbox/node/hub.go b/services/agents-api/internal/sandbox/node/hub.go index 88284fece..5a1b4d89e 100644 --- a/services/agents-api/internal/sandbox/node/hub.go +++ b/services/agents-api/internal/sandbox/node/hub.go @@ -380,7 +380,7 @@ func (h *Hub) call(ctx context.Context, id string, q request) (response, error) defer timer.Stop() select { case result := <-ch: - return result, responseError(result.ErrorCode) + return result, responseError(result) case <-p.done: return response{}, uncertain(q.Operation, ErrUnavailable) case <-ctx.Done(): diff --git a/services/agents-api/internal/sandbox/node/node_test.go b/services/agents-api/internal/sandbox/node/node_test.go index eed327bf7..6e1e5f424 100644 --- a/services/agents-api/internal/sandbox/node/node_test.go +++ b/services/agents-api/internal/sandbox/node/node_test.go @@ -166,7 +166,7 @@ func TestLostCreateResponseDoesNotReplayAndReconnectSerializesCleanup(t *testing func TestOfflineIsUnknownAndDockerDoesNotAdvertiseCheckpoint(t *testing.T) { h := NewHub(HubOptions{OwnerEpoch: func(context.Context) (uint64, error) { return 1, nil }}) p := h.Proxy(uuid.NewString(), "docker", 1) - if _, ok := p.(sandbox.CheckpointProvider); ok { + if sandbox.SupportsCheckpoint(p) { t.Fatal("docker advertised checkpoint") } _, err := p.GetInfo(context.Background(), reference()) diff --git a/services/agents-api/internal/sandbox/node/observations.go b/services/agents-api/internal/sandbox/node/observations.go index c63c3215c..dc731fc5d 100644 --- a/services/agents-api/internal/sandbox/node/observations.go +++ b/services/agents-api/internal/sandbox/node/observations.go @@ -3,6 +3,7 @@ package node import ( "context" "errors" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeobs" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" @@ -40,9 +41,12 @@ func (p *provider) Observe(ctx context.Context, target runtimeobs.Target) (runti } func observeProvider(ctx context.Context, provider sandbox.SandboxProvider, target runtimeobs.Target) (runtimeobs.Sample, error) { + if err := providercontract.Require(provider, "Observe"); err != nil { + return runtimeobs.Sample{}, err + } source, ok := provider.(runtimeobs.Source) if !ok { - return runtimeobs.Sample{}, runtimeobs.ErrUnavailable + return runtimeobs.Sample{}, providercontract.ErrContract } return source.Observe(ctx, target) } diff --git a/services/agents-api/internal/sandbox/node/observations_test.go b/services/agents-api/internal/sandbox/node/observations_test.go index 952227ebe..c41ab6a2a 100644 --- a/services/agents-api/internal/sandbox/node/observations_test.go +++ b/services/agents-api/internal/sandbox/node/observations_test.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "errors" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" "net/http/httptest" "reflect" "sync" @@ -143,7 +144,7 @@ func TestObservationWirePreservesUnavailableAndRejectsMismatchedIdentity(t *test provider sandbox.SandboxProvider want error }{ - {"unsupported", &fakeProvider{}, runtimeobs.ErrUnavailable}, + {"unsupported", &fakeProvider{}, providercontract.ErrUnsupported}, {"unavailable", &observationProvider{fakeProvider: &fakeProvider{}, err: runtimeobs.ErrUnavailable}, runtimeobs.ErrUnavailable}, {"stopped", &observationProvider{fakeProvider: &fakeProvider{}, err: runtimeobs.ErrNotRunning}, runtimeobs.ErrNotRunning}, {"ownership", &observationProvider{fakeProvider: &fakeProvider{}, err: sandbox.ErrOwnership}, sandbox.ErrOwnership}, @@ -151,7 +152,7 @@ func TestObservationWirePreservesUnavailableAndRejectsMismatchedIdentity(t *test for _, test := range providers { t.Run(test.name, func(t *testing.T) { out := execute(t.Context(), test.provider, q) - if !errors.Is(responseError(out.ErrorCode), test.want) || out.Sample != nil { + if !errors.Is(responseError(out), test.want) || out.Sample != nil { t.Fatal("observation error or missing sample was fabricated", out) } }) @@ -161,7 +162,7 @@ func TestObservationWirePreservesUnavailableAndRejectsMismatchedIdentity(t *test invalid := target mutate(&invalid) q.Observation = &invalid - if out := execute(t.Context(), p, q); !errors.Is(responseError(out.ErrorCode), sandbox.ErrInvalid) { + if out := execute(t.Context(), p, q); !errors.Is(responseError(out), sandbox.ErrInvalid) { t.Fatal("mismatched observation reached provider", out) } } diff --git a/services/agents-api/internal/sandbox/node/operations.go b/services/agents-api/internal/sandbox/node/operations.go new file mode 100644 index 000000000..8f016664b --- /dev/null +++ b/services/agents-api/internal/sandbox/node/operations.go @@ -0,0 +1,30 @@ +package node + +// Transport operation names map to the authored Provider method declarations. +var operationMethods = map[string]string{ + "create": "Create", + "info": "GetInfo", + "renew": "Renew", + "kill": "Kill", + "command": "RunCommand", + "initial": "Initial", + "new_compute": "NewCompute", + "compute": "GetCompute", + "suspend": "Suspend", + "resume": "Resume", + "kill_compute": "KillCompute", + "delete_snapshot": "DeleteSnapshot", + "command_compute": "RunCommandCompute", + "resume_compute": "ResumeCompute", + "observe": "Observe", +} + +func operationMethod(wire string) string { return operationMethods[wire] } +func operationWire(method string) string { + for wire, name := range operationMethods { + if name == method { + return wire + } + } + return "" +} diff --git a/services/agents-api/internal/sandbox/node/operations_test.go b/services/agents-api/internal/sandbox/node/operations_test.go new file mode 100644 index 000000000..67bbd6292 --- /dev/null +++ b/services/agents-api/internal/sandbox/node/operations_test.go @@ -0,0 +1,58 @@ +package node + +import ( + "context" + "encoding/json" + "errors" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" + "github.com/google/uuid" + "testing" +) + +func TestUnsupportedWireIsExplicitAndDoesNotInvokeProvider(t *testing.T) { + p := &fakeProvider{} + q := request{ID: uuid.NewString(), ConnectionID: uuid.NewString(), Reference: reference(), TimeoutMillis: 1000, Operation: "initial"} + out := execute(t.Context(), p, q) + if p.creates != 0 || p.reads != 0 || p.kills != 0 { + t.Fatal("unsupported operation reached provider") + } + raw, err := json.Marshal(out) + if err != nil { + t.Fatal(err) + } + var decoded response + if err := json.Unmarshal(raw, &decoded); err != nil { + t.Fatal(err) + } + if reason, ok := providercontract.UnsupportedReason(responseError(decoded), "Initial"); !ok || reason != "fixture_operation_not_supported" { + t.Fatal(decoded) + } + for _, bad := range []response{{ErrorCode: "unsupported"}, {ErrorCode: "unsupported", Unsupported: &providercontract.UnsupportedError{Operation: "Initial", Reason: "private / credential"}}, {Unsupported: out.Unsupported}} { + if !errors.Is(responseError(bad), sandbox.ErrComputeUnconfirmed) { + t.Fatal("malformed failure became certainty", bad) + } + } +} +func TestUnsupportedProxyRejectsBeforeNodeResolution(t *testing.T) { + p := &provider{kind: "docker", resolveGeneration: func(context.Context, sandbox.Reference) (string, uint64, error) { + t.Fatal("unsupported call resolved a node") + return "", 0, nil + }} + if err := sandbox.ValidateProvider(p); err != nil { + t.Fatal(err) + } + if _, err := p.Initial(t.Context(), reference()); !errors.Is(err, providercontract.ErrUnsupported) { + t.Fatal(err) + } + if _, err := p.ObserveBatch(t.Context(), nil); !errors.Is(err, providercontract.ErrUnsupported) { + t.Fatal(err) + } +} +func TestNodeOperationMappingCoversForwardedMethods(t *testing.T) { + for _, method := range []string{"Create", "GetInfo", "Renew", "Kill", "RunCommand", "Initial", "NewCompute", "GetCompute", "Suspend", "Resume", "KillCompute", "DeleteSnapshot", "RunCommandCompute", "ResumeCompute", "Observe"} { + if wire := operationWire(method); wire == "" || operationMethod(wire) != method { + t.Fatal(method) + } + } +} diff --git a/services/agents-api/internal/sandbox/node/provider_operations_fixture_test.go b/services/agents-api/internal/sandbox/node/provider_operations_fixture_test.go new file mode 100644 index 000000000..5a8b8d579 --- /dev/null +++ b/services/agents-api/internal/sandbox/node/provider_operations_fixture_test.go @@ -0,0 +1,92 @@ +package node + +import ( + "context" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeobs" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" +) + +func (*fakeProvider) ProviderOperations() providercontract.Operations { + return providercontract.Operations{ + "Create": {State: providercontract.Supported}, + "GetInfo": {State: providercontract.Supported}, + "Renew": {State: providercontract.Supported}, + "Kill": {State: providercontract.Supported}, + "RunCommand": {State: providercontract.Supported}, + "Initial": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "NewCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "GetCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "Suspend": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "Resume": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "KillCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "DeleteSnapshot": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "RunCommandCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "ResumeCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "Observe": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "ObserveBatch": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "DiscoverSelection": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "VerifyCredential": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + } +} +func (*fakeProvider) Initial(context.Context, sandbox.Reference) (sandbox.Compute, error) { + return sandbox.Compute{}, &providercontract.UnsupportedError{Operation: "Initial", Reason: "fixture_operation_not_supported"} +} +func (*fakeProvider) NewCompute(context.Context, sandbox.Reference, uint64, *sandbox.SnapshotIdentity) (sandbox.Compute, error) { + return sandbox.Compute{}, &providercontract.UnsupportedError{Operation: "NewCompute", Reason: "fixture_operation_not_supported"} +} +func (*fakeProvider) GetCompute(context.Context, sandbox.Reference, sandbox.Compute) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "GetCompute", Reason: "fixture_operation_not_supported"} +} +func (*fakeProvider) Suspend(context.Context, sandbox.SuspendRequest) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "Suspend", Reason: "fixture_operation_not_supported"} +} +func (*fakeProvider) Resume(context.Context, sandbox.ResumeRequest) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "Resume", Reason: "fixture_operation_not_supported"} +} +func (*fakeProvider) KillCompute(context.Context, sandbox.Reference, sandbox.Compute) error { + return &providercontract.UnsupportedError{Operation: "KillCompute", Reason: "fixture_operation_not_supported"} +} +func (*fakeProvider) DeleteSnapshot(context.Context, sandbox.Reference, sandbox.SnapshotIdentity) error { + return &providercontract.UnsupportedError{Operation: "DeleteSnapshot", Reason: "fixture_operation_not_supported"} +} +func (*fakeProvider) RunCommandCompute(context.Context, sandbox.Reference, sandbox.Compute, sandbox.Command) (sandbox.CommandResult, error) { + return sandbox.CommandResult{}, &providercontract.UnsupportedError{Operation: "RunCommandCompute", Reason: "fixture_operation_not_supported"} +} +func (*fakeProvider) ResumeCompute(context.Context, sandbox.Reference, sandbox.Compute) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "ResumeCompute", Reason: "fixture_operation_not_supported"} +} +func (*fakeProvider) Observe(context.Context, runtimeobs.Target) (runtimeobs.Sample, error) { + return runtimeobs.Sample{}, &providercontract.UnsupportedError{Operation: "Observe", Reason: "fixture_operation_not_supported"} +} +func (*fakeProvider) ObserveBatch(context.Context, []runtimeobs.Target) ([]runtimeobs.BatchResult, error) { + return nil, &providercontract.UnsupportedError{Operation: "ObserveBatch", Reason: "fixture_operation_not_supported"} +} +func (*fakeProvider) DiscoverSelection(context.Context, sandbox.Selection) (sandbox.Selection, error) { + return sandbox.Selection{}, &providercontract.UnsupportedError{Operation: "DiscoverSelection", Reason: "fixture_operation_not_supported"} +} +func (*fakeProvider) VerifyCredential(context.Context, []sandbox.Reference) error { + return &providercontract.UnsupportedError{Operation: "VerifyCredential", Reason: "fixture_operation_not_supported"} +} +func (*observationProvider) ProviderOperations() providercontract.Operations { + return providercontract.Operations{ + "Create": {State: providercontract.Supported}, + "GetInfo": {State: providercontract.Supported}, + "Renew": {State: providercontract.Supported}, + "Kill": {State: providercontract.Supported}, + "RunCommand": {State: providercontract.Supported}, + "Initial": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "NewCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "GetCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "Suspend": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "Resume": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "KillCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "DeleteSnapshot": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "RunCommandCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "ResumeCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "Observe": {State: providercontract.Supported}, + "ObserveBatch": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "DiscoverSelection": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "VerifyCredential": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + } +} diff --git a/services/agents-api/internal/sandbox/node/proxy.go b/services/agents-api/internal/sandbox/node/proxy.go index 1b921c190..ec0d7ba77 100644 --- a/services/agents-api/internal/sandbox/node/proxy.go +++ b/services/agents-api/internal/sandbox/node/proxy.go @@ -2,6 +2,9 @@ package node import ( "context" + "errors" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeobs" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox/providers" @@ -12,10 +15,9 @@ type provider struct { resolveGeneration func(context.Context, sandbox.Reference) (string, uint64, error) kind string } -type checkpointProvider struct{ *provider } var _ sandbox.SandboxProvider = (*provider)(nil) -var _ sandbox.CheckpointProvider = (*checkpointProvider)(nil) +var _ sandbox.CheckpointProvider = (*provider)(nil) // Proxy binds a fixed node and deployment generation explicitly. func (h *Hub) Proxy(id, kind string, generation uint64) sandbox.SandboxProvider { @@ -23,6 +25,9 @@ func (h *Hub) Proxy(id, kind string, generation uint64) sandbox.SandboxProvider } func (p *provider) call(ctx context.Context, q request) (response, error) { + if err := providercontract.Require(p, operationMethod(q.Operation)); err != nil { + return response{}, err + } id, generation, err := p.resolveGeneration(ctx, q.Reference) if err != nil { return response{}, err @@ -31,7 +36,13 @@ func (p *provider) call(ctx context.Context, q request) (response, error) { return response{}, sandbox.ErrOwnership } q.DeploymentGeneration = generation - return p.hub.call(ctx, id, q) + out, err := p.hub.call(ctx, id, q) + if errors.Is(err, providercontract.ErrUnsupported) { + if _, valid := providercontract.UnsupportedReason(err, operationMethod(q.Operation)); !valid { + return response{}, sandbox.ErrComputeUnconfirmed + } + } + return out, err } func (p *provider) info(ctx context.Context, q request) (sandbox.Info, error) { r, e := p.call(ctx, q) @@ -58,6 +69,26 @@ func creationSettled(info *sandbox.Info, ref sandbox.Reference) bool { } return info.ProviderID != "" } +func (p *provider) ProviderOperations() providercontract.Operations { + adapter, err := providers.Lookup(p.kind) + if err != nil { + return nil + } + ops := adapter.Operations() + 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"} + return ops +} +func (*provider) ObserveBatch(context.Context, []runtimeobs.Target) ([]runtimeobs.BatchResult, error) { + return nil, &providercontract.UnsupportedError{Operation: "ObserveBatch", Reason: "node_transport_has_no_batch_observation"} +} +func (p *provider) DiscoverSelection(context.Context, sandbox.Selection) (sandbox.Selection, error) { + return sandbox.Selection{}, &providercontract.UnsupportedError{Operation: "DiscoverSelection", Reason: "node_configuration_is_core_owned"} +} +func (p *provider) VerifyCredential(context.Context, []sandbox.Reference) error { + return &providercontract.UnsupportedError{Operation: "VerifyCredential", Reason: "node_credentials_are_transport_owned"} +} func (p *provider) Create(ctx context.Context, b sandbox.Bootstrap) (sandbox.Info, error) { return p.info(ctx, request{Operation: "create", Reference: b.Reference, Bootstrap: &b}) } @@ -84,7 +115,7 @@ func (p *provider) command(ctx context.Context, q request) (sandbox.CommandResul func (p *provider) RunCommand(ctx context.Context, r sandbox.Reference, c sandbox.Command) (sandbox.CommandResult, error) { return p.command(ctx, request{Operation: "command", Reference: r, Command: &c}) } -func (p *checkpointProvider) Initial(ctx context.Context, r sandbox.Reference) (sandbox.Compute, error) { +func (p *provider) Initial(ctx context.Context, r sandbox.Reference) (sandbox.Compute, error) { out, e := p.call(ctx, request{Operation: "initial", Reference: r}) if e != nil { return sandbox.Compute{}, e @@ -94,7 +125,7 @@ func (p *checkpointProvider) Initial(ctx context.Context, r sandbox.Reference) ( } return *out.Compute, nil } -func (p *checkpointProvider) NewCompute(ctx context.Context, r sandbox.Reference, g uint64, s *sandbox.SnapshotIdentity) (sandbox.Compute, error) { +func (p *provider) NewCompute(ctx context.Context, r sandbox.Reference, g uint64, s *sandbox.SnapshotIdentity) (sandbox.Compute, error) { out, e := p.call(ctx, request{Operation: "new_compute", Reference: r, Generation: g, Snapshot: s}) if e != nil { return sandbox.Compute{}, e @@ -104,7 +135,7 @@ func (p *checkpointProvider) NewCompute(ctx context.Context, r sandbox.Reference } return *out.Compute, nil } -func (p *checkpointProvider) state(ctx context.Context, q request) (sandbox.ComputeState, error) { +func (p *provider) state(ctx context.Context, q request) (sandbox.ComputeState, error) { r, e := p.call(ctx, q) if e != nil { return sandbox.ComputeState{}, e @@ -114,27 +145,27 @@ func (p *checkpointProvider) state(ctx context.Context, q request) (sandbox.Comp } return *r.State, nil } -func (p *checkpointProvider) GetCompute(ctx context.Context, r sandbox.Reference, c sandbox.Compute) (sandbox.ComputeState, error) { +func (p *provider) GetCompute(ctx context.Context, r sandbox.Reference, c sandbox.Compute) (sandbox.ComputeState, error) { return p.state(ctx, request{Operation: "compute", Reference: r, Compute: &c}) } -func (p *checkpointProvider) Suspend(ctx context.Context, q sandbox.SuspendRequest) (sandbox.ComputeState, error) { +func (p *provider) Suspend(ctx context.Context, q sandbox.SuspendRequest) (sandbox.ComputeState, error) { return p.state(ctx, request{Operation: "suspend", Reference: q.Reference, Suspend: &q}) } -func (p *checkpointProvider) Resume(ctx context.Context, q sandbox.ResumeRequest) (sandbox.ComputeState, error) { +func (p *provider) Resume(ctx context.Context, q sandbox.ResumeRequest) (sandbox.ComputeState, error) { return p.state(ctx, request{Operation: "resume", Reference: q.Reference, Resume: &q}) } -func (p *checkpointProvider) KillCompute(ctx context.Context, r sandbox.Reference, c sandbox.Compute) error { +func (p *provider) KillCompute(ctx context.Context, r sandbox.Reference, c sandbox.Compute) error { _, e := p.call(ctx, request{Operation: "kill_compute", Reference: r, Compute: &c}) return e } -func (p *checkpointProvider) DeleteSnapshot(ctx context.Context, r sandbox.Reference, s sandbox.SnapshotIdentity) error { +func (p *provider) DeleteSnapshot(ctx context.Context, r sandbox.Reference, s sandbox.SnapshotIdentity) error { _, e := p.call(ctx, request{Operation: "delete_snapshot", Reference: r, Snapshot: &s}) return e } -func (p *checkpointProvider) RunCommandCompute(ctx context.Context, r sandbox.Reference, c sandbox.Compute, v sandbox.Command) (sandbox.CommandResult, error) { +func (p *provider) RunCommandCompute(ctx context.Context, r sandbox.Reference, c sandbox.Compute, v sandbox.Command) (sandbox.CommandResult, error) { return p.command(ctx, request{Operation: "command_compute", Reference: r, Compute: &c, Command: &v}) } -func (p *checkpointProvider) ResumeCompute(ctx context.Context, r sandbox.Reference, c sandbox.Compute) (sandbox.ComputeState, error) { +func (p *provider) ResumeCompute(ctx context.Context, r sandbox.Reference, c sandbox.Compute) (sandbox.ComputeState, error) { return p.state(ctx, request{Operation: "resume_compute", Reference: r, Compute: &c}) } @@ -142,8 +173,5 @@ func (p *checkpointProvider) ResumeCompute(ctx context.Context, r sandbox.Refere // distinct from the request's compute generation. func (h *Hub) GenerationProvider(kind string, resolve func(context.Context, sandbox.Reference) (string, uint64, error)) sandbox.SandboxProvider { p := &provider{hub: h, kind: kind, resolveGeneration: resolve} - if providers.SupportsCheckpoint(kind) { - return &checkpointProvider{p} - } return p } diff --git a/services/agents-api/internal/sandbox/node/wire.go b/services/agents-api/internal/sandbox/node/wire.go index c735b96fd..53a34f2c2 100644 --- a/services/agents-api/internal/sandbox/node/wire.go +++ b/services/agents-api/internal/sandbox/node/wire.go @@ -7,6 +7,7 @@ import ( "context" "encoding/json" "errors" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" "io" "time" @@ -16,7 +17,7 @@ import ( "github.com/gorilla/websocket" ) -const ProtocolVersion = 3 +const ProtocolVersion = 4 const MaxControlFrameBytes = 32 * 1024 const MaxFrameBytes = 72 * 1024 * 1024 const maxPending = 32 @@ -96,14 +97,15 @@ type request struct { } type response struct { - ID string `json:"id"` - ConnectionID string `json:"connection_id"` - ErrorCode string `json:"error_code,omitempty"` - Info *sandbox.Info `json:"info,omitempty"` - Compute *sandbox.Compute `json:"compute,omitempty"` - State *sandbox.ComputeState `json:"state,omitempty"` - Command *sandbox.CommandResult `json:"command,omitempty"` - Sample *runtimeobs.Sample `json:"sample,omitempty"` + Unsupported *providercontract.UnsupportedError `json:"unsupported,omitempty"` + ID string `json:"id"` + ConnectionID string `json:"connection_id"` + ErrorCode string `json:"error_code,omitempty"` + Info *sandbox.Info `json:"info,omitempty"` + Compute *sandbox.Compute `json:"compute,omitempty"` + State *sandbox.ComputeState `json:"state,omitempty"` + Command *sandbox.CommandResult `json:"command,omitempty"` + Sample *runtimeobs.Sample `json:"sample,omitempty"` } type frame struct { @@ -171,6 +173,8 @@ func errorCode(err error) string { switch { case err == nil: return "" + case errors.Is(err, providercontract.ErrUnsupported): + return "unsupported" case errors.Is(err, runtimeobs.ErrUnavailable): return "observation_unavailable" case errors.Is(err, runtimeobs.ErrNotRunning): @@ -189,8 +193,20 @@ func errorCode(err error) string { return "unconfirmed" } } -func responseError(code string) error { - switch code { +func responseError(out response) error { + if out.ErrorCode == "unsupported" { + if out.Unsupported == nil || operationWire(out.Unsupported.Operation) == "" { + return sandbox.ErrComputeUnconfirmed + } + if _, valid := providercontract.UnsupportedReason(out.Unsupported, out.Unsupported.Operation); !valid { + return sandbox.ErrComputeUnconfirmed + } + return out.Unsupported + } + if out.Unsupported != nil { + return sandbox.ErrComputeUnconfirmed + } + switch out.ErrorCode { case "": return nil case "observation_unavailable": @@ -278,6 +294,14 @@ func execute(ctx context.Context, p sandbox.SandboxProvider, q request) response out.ErrorCode = errorCode(err) return out } + if err := providercontract.Require(p, operationMethod(q.Operation)); err != nil { + out.ErrorCode = errorCode(err) + var unsupported *providercontract.UnsupportedError + if errors.As(err, &unsupported) { + out.Unsupported = unsupported + } + return out + } var info sandbox.Info var command sandbox.CommandResult switch q.Operation { @@ -300,9 +324,9 @@ func execute(ctx context.Context, p sandbox.SandboxProvider, q request) response command, err = p.RunCommand(ctx, q.Reference, *q.Command) out.Command = &command default: - cp, ok := p.(sandbox.CheckpointProvider) - if !ok { - err = sandbox.ErrInvalid + cp, checkpointErr := sandbox.Checkpoint(p) + if checkpointErr != nil { + err = checkpointErr break } var state sandbox.ComputeState @@ -338,6 +362,10 @@ func execute(ctx context.Context, p sandbox.SandboxProvider, q request) response } } out.ErrorCode = errorCode(err) + var unsupported *providercontract.UnsupportedError + if errors.As(err, &unsupported) { + out.Unsupported = unsupported + } if err != nil { out.Sample = nil if !creationSettled(out.Info, q.Reference) { diff --git a/services/agents-api/internal/sandbox/operations.go b/services/agents-api/internal/sandbox/operations.go new file mode 100644 index 000000000..6c781cb62 --- /dev/null +++ b/services/agents-api/internal/sandbox/operations.go @@ -0,0 +1,80 @@ +package sandbox + +import ( + "errors" + "fmt" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeobs" + "reflect" +) + +// 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[SelectionDiscoverer](), reflect.TypeFor[CredentialVerifier](), + reflect.TypeFor[runtimeobs.Source](), reflect.TypeFor[runtimeobs.BatchSource](), +} + +func ValidateProvider(p SandboxProvider) error { + if err := providercontract.Validate(p, providerInterfaces...); err != nil { + return err + } + return ValidateOperations(p.ProviderOperations()) +} + +// ValidateOperations rejects omitted, unknown and contradictory declarations. +func ValidateOperations(operations providercontract.Operations) error { + known := map[string]bool{} + for _, contract := range providerInterfaces { + for i := 0; i < contract.NumMethod(); i++ { + name := contract.Method(i).Name + if name != "ProviderOperations" { + known[name] = true + if err := operations[name].Check(name); err != nil && !errors.Is(err, providercontract.ErrUnsupported) { + return err + } + } + } + } + for name := range operations { + if !known[name] { + return fmt.Errorf("%w: unknown operation %s", providercontract.ErrContract, name) + } + } + required := reflect.TypeFor[SandboxProvider]() + for i := 0; i < required.NumMethod(); i++ { + name := required.Method(i).Name + if name != "ProviderOperations" && operations[name].State != providercontract.Supported { + return fmt.Errorf("%w: required operation %s", providercontract.ErrContract, name) + } + } + // The checkpoint lifecycle is indivisible: partial cleanup or restore support + // cannot safely own a compute incarnation. + checkpoint := reflect.TypeFor[CheckpointProvider]() + for i := 0; i < checkpoint.NumMethod(); i++ { + name := checkpoint.Method(i).Name + if _, required := reflect.TypeFor[SandboxProvider]().MethodByName(name); !required && operations[name].State != operations["Initial"].State { + return fmt.Errorf("%w: incomplete checkpoint lifecycle", providercontract.ErrContract) + } + } + if operations["ObserveBatch"].State == providercontract.Supported && operations["Observe"].State != providercontract.Supported { + return fmt.Errorf("%w: batch observation requires observation", providercontract.ErrContract) + } + return nil +} + +func SupportsCheckpoint(p SandboxProvider) bool { + return providercontract.Require(p, "Initial") == nil +} + +func Checkpoint(p SandboxProvider) (CheckpointProvider, error) { + if err := providercontract.Require(p, "Initial"); err != nil { + return nil, err + } + cp, ok := p.(CheckpointProvider) + if !ok { + return nil, providercontract.ErrContract + } + return cp, nil +} diff --git a/services/agents-api/internal/sandbox/operations_test.go b/services/agents-api/internal/sandbox/operations_test.go new file mode 100644 index 000000000..93cb1320a --- /dev/null +++ b/services/agents-api/internal/sandbox/operations_test.go @@ -0,0 +1,112 @@ +package sandbox_test + +import ( + "context" + "errors" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox/docker" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox/e2b" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox/microsandbox" + "reflect" + "testing" +) + +type changedDeclaration struct { + *docker.Provider + operations providercontract.Operations +} + +func (p *changedDeclaration) ProviderOperations() providercontract.Operations { return p.operations } + +type onlyRequired struct{ sandbox.SandboxProvider } + +func (*onlyRequired) ProviderOperations() providercontract.Operations { return docker.Operations() } + +func TestDeclarationsRejectMissingUnknownAndContradictoryOperations(t *testing.T) { + for _, mutate := range []struct { + name string + change func(providercontract.Operations) + }{ + {"omitted", func(o providercontract.Operations) { delete(o, "DeleteSnapshot") }}, + {"zero", func(o providercontract.Operations) { o["ObserveBatch"] = providercontract.Support{} }}, + {"unknown", func(o providercontract.Operations) { + o["FutureOperation"] = providercontract.Support{State: providercontract.Supported} + }}, + {"required unsupported", func(o providercontract.Operations) { + o["Renew"] = providercontract.Support{State: providercontract.Unsupported, Reason: "no_native_lease"} + }}, + {"partial checkpoint", func(o providercontract.Operations) { + o["Initial"] = providercontract.Support{State: providercontract.Supported} + }}, + {"unsafe reason", func(o providercontract.Operations) { + o["ObserveBatch"] = providercontract.Support{State: providercontract.Unsupported, Reason: "https://private:key@host"} + }}, + {"unknown state", func(o providercontract.Operations) { + o["ObserveBatch"] = providercontract.Support{State: "unavailable"} + }}, + } { + t.Run(mutate.name, func(t *testing.T) { + o := docker.Operations() + mutate.change(o) + if err := sandbox.ValidateProvider(&changedDeclaration{Provider: &docker.Provider{}, operations: o}); !errors.Is(err, providercontract.ErrContract) { + t.Fatal(err) + } + }) + } + if err := sandbox.ValidateProvider(&onlyRequired{}); !errors.Is(err, providercontract.ErrContract) { + t.Fatal("missing extension implementations accepted", err) + } + var nilProvider *docker.Provider + if err := sandbox.ValidateProvider(nilProvider); !errors.Is(err, providercontract.ErrContract) { + t.Fatal("typed nil accepted", err) + } +} + +// Unsupported methods must be callable safely even without native clients or +// configuration: any attempt at native I/O would panic on these zero providers. +func TestEveryUnsupportedNativeOperationRejectsWithoutSideEffects(t *testing.T) { + for _, provider := range []sandbox.SandboxProvider{&docker.Provider{}, &e2b.Provider{}, µsandbox.Provider{}} { + if err := sandbox.ValidateProvider(provider); err != nil { + t.Fatal(err) + } + for name, support := range provider.ProviderOperations() { + if support.State != providercontract.Unsupported { + continue + } + method := reflect.ValueOf(provider).MethodByName(name) + arguments := make([]reflect.Value, method.Type().NumIn()) + for i := range arguments { + arguments[i] = reflect.Zero(method.Type().In(i)) + } + arguments[0] = reflect.ValueOf(context.Background()) + values := method.Call(arguments) + err, _ := values[len(values)-1].Interface().(error) + reason, ok := providercontract.UnsupportedReason(err, name) + if !ok || reason != support.Reason { + t.Fatalf("%T.%s returned %v", provider, name, err) + } + for _, v := range values[:len(values)-1] { + if !v.IsZero() { + t.Fatalf("%s fabricated a successful value", name) + } + } + } + } +} + +// An extended interface cannot inherit success through the existing declaration. +type nextContract interface { + sandbox.SandboxProvider + NextOperation(context.Context) error +} +type futureProvider struct{ *docker.Provider } + +func (*futureProvider) NextOperation(context.Context) error { return nil } +func TestNewContractRequiresAnAuthoredDecision(t *testing.T) { + for _, p := range []providercontract.Declared{&docker.Provider{}, &futureProvider{&docker.Provider{}}} { + if err := providercontract.Validate(p, reflect.TypeFor[nextContract]()); !errors.Is(err, providercontract.ErrContract) { + t.Fatal("new operation inherited a default", err) + } + } +} diff --git a/services/agents-api/internal/sandbox/providers/config.go b/services/agents-api/internal/sandbox/providers/config.go index 61839d052..8e71e4a46 100644 --- a/services/agents-api/internal/sandbox/providers/config.go +++ b/services/agents-api/internal/sandbox/providers/config.go @@ -83,6 +83,10 @@ func Build(config Config) (*Built, func(), error) { if err != nil { return nil, closeProvider, err } + if err := ValidateBinding(adapter, result.Provider); err != nil { + closeProvider() + return nil, func() {}, err + } return result, closeProvider, nil } diff --git a/services/agents-api/internal/sandbox/providers/e2b.go b/services/agents-api/internal/sandbox/providers/e2b.go index 783f1a47f..8f7b4f54c 100644 --- a/services/agents-api/internal/sandbox/providers/e2b.go +++ b/services/agents-api/internal/sandbox/providers/e2b.go @@ -22,7 +22,14 @@ func BuildDirect(c DirectConfig) (sandbox.SandboxProvider, error) { if a.BuildDirect == nil { return nil, sandbox.ErrInvalid } - return a.BuildDirect(c) + p, err := a.BuildDirect(c) + if err != nil { + return nil, err + } + if err := ValidateBinding(a, p); err != nil { + return nil, err + } + return p, nil } func buildE2B(c DirectConfig) (sandbox.SandboxProvider, error) { if c.Selection.E2B == nil { diff --git a/services/agents-api/internal/sandbox/providers/operations.go b/services/agents-api/internal/sandbox/providers/operations.go new file mode 100644 index 000000000..b6c619c5d --- /dev/null +++ b/services/agents-api/internal/sandbox/providers/operations.go @@ -0,0 +1,18 @@ +package providers + +import ( + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" + "reflect" +) + +// ValidateBinding catches construction that disagrees with its registration. +func ValidateBinding(adapter Adapter, provider sandbox.SandboxProvider) error { + if err := sandbox.ValidateProvider(provider); err != nil { + return err + } + if adapter.Operations == nil || !reflect.DeepEqual(adapter.Operations(), provider.ProviderOperations()) { + return providercontract.ErrContract + } + return nil +} diff --git a/services/agents-api/internal/sandbox/providers/operations_test.go b/services/agents-api/internal/sandbox/providers/operations_test.go new file mode 100644 index 000000000..ad305f48a --- /dev/null +++ b/services/agents-api/internal/sandbox/providers/operations_test.go @@ -0,0 +1,22 @@ +package providers + +import ( + "errors" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox/docker" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox/e2b" + "testing" +) + +func TestRegistrationRejectsMissingAndMismatchedDeclarations(t *testing.T) { + for _, adapter := range []Adapter{{}, {Operations: func() providercontract.Operations { return nil }}} { + adapters["invalid-contract-fixture"] = adapter + if _, err := Lookup("invalid-contract-fixture"); err == nil { + t.Fatal("invalid declaration registered") + } + } + delete(adapters, "invalid-contract-fixture") + if err := ValidateBinding(Adapter{Operations: e2b.Operations}, &docker.Provider{}); !errors.Is(err, providercontract.ErrContract) { + t.Fatal("registration differs from instance", err) + } +} diff --git a/services/agents-api/internal/sandbox/providers/registry.go b/services/agents-api/internal/sandbox/providers/registry.go index a06a2c30e..3591fc145 100644 --- a/services/agents-api/internal/sandbox/providers/registry.go +++ b/services/agents-api/internal/sandbox/providers/registry.go @@ -6,6 +6,7 @@ import ( "crypto/sha256" "encoding/hex" "fmt" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox/docker" @@ -14,7 +15,7 @@ import ( ) // Adapter describes configuration and transport independently of compute operations. -// Optional native operations remain interface assertions on SandboxProvider. +// Native operation support comes from the adapter-owned complete declaration. type Adapter struct { Policy sandbox.DeploymentPolicy ReplaceCredential func(sandbox.Selection, sandbox.Selection) sandbox.Selection @@ -27,7 +28,7 @@ type Adapter struct { Credential bool PublicOrigin bool Mode string - Checkpoint bool + Operations func() providercontract.Operations IdleSeconds, RetentionSeconds int64 ValidateSpecification func(sandbox.DeploymentSpec) error ValidateResources func(sandbox.Resources) error @@ -36,18 +37,18 @@ type Adapter struct { var adapters = map[string]Adapter{ "docker": { - Policy: docker.Policy(), Mode: "nodes", BuildLocal: buildDocker, + Policy: docker.Policy(), Operations: docker.Operations, Mode: "nodes", BuildLocal: buildDocker, ValidateSpecification: docker.ValidateSpecification, ValidateResources: docker.ValidateResources, Normalize: nodeSelection(docker.ValidateSpecification), }, "microsandbox": { - Policy: microsandbox.Policy(), Mode: "nodes", BuildLocal: buildMicrosandbox, - Checkpoint: true, IdleSeconds: 300, RetentionSeconds: 86400, + Policy: microsandbox.Policy(), Operations: microsandbox.Operations, Mode: "nodes", BuildLocal: buildMicrosandbox, + IdleSeconds: 300, RetentionSeconds: 86400, ValidateSpecification: microsandbox.ValidateSpecification, ValidateResources: microsandbox.ValidateResources, Normalize: nodeSelection(microsandbox.ValidateSpecification), }, "e2b": { - Policy: e2b.Policy(), Mode: "direct", BuildDirect: buildE2B, + Policy: e2b.Policy(), Operations: e2b.Operations, Mode: "direct", BuildDirect: buildE2B, Credential: true, PublicOrigin: true, ReplaceCredential: e2b.ReplaceCredential, CredentialRequiresReset: e2b.CredentialRequiresReset, CredentialUnconfirmed: e2b.ErrRequestUnconfirmed, Restore: e2b.RestoreSelection, @@ -58,18 +59,24 @@ var adapters = map[string]Adapter{ func Lookup(kind string) (Adapter, error) { a, ok := adapters[kind] - if !ok { + if !ok || a.Operations == nil { return Adapter{}, fmt.Errorf("%w: unsupported sandbox provider", sandbox.ErrInvalid) } + if err := sandbox.ValidateOperations(a.Operations()); err != nil { + return Adapter{}, err + } return a, nil } -func IsNode(kind string) bool { a, e := Lookup(kind); return e == nil && a.Mode == "nodes" } -func SupportsCheckpoint(kind string) bool { a, e := Lookup(kind); return e == nil && a.Checkpoint } +func IsNode(kind string) bool { a, e := Lookup(kind); return e == nil && a.Mode == "nodes" } +func SupportsCheckpoint(kind string) bool { + a, e := Lookup(kind) + return e == nil && a.Operations()["Initial"].State == providercontract.Supported +} // RetainedLimit keeps nodes without checkpoint support within their active capacity. func RetainedLimit(kind string, active, retained int) int { a, err := Lookup(kind) - if err == nil && a.Mode == "nodes" && !a.Checkpoint { + if err == nil && a.Mode == "nodes" && !SupportsCheckpoint(kind) { return active } return retained diff --git a/services/agents-api/internal/sandbox/providers/registry_test.go b/services/agents-api/internal/sandbox/providers/registry_test.go index 17212d63a..51b99954b 100644 --- a/services/agents-api/internal/sandbox/providers/registry_test.go +++ b/services/agents-api/internal/sandbox/providers/registry_test.go @@ -2,6 +2,8 @@ package providers import ( "errors" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox/docker" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox/microsandbox" "testing" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" @@ -63,7 +65,7 @@ func TestSelectionNormalizationAndCredentialInheritance(t *testing.T) { func TestNewRegistrationDoesNotNeedCoreDispatchChanges(t *testing.T) { const kind = "contract-test-provider" // Registration is test-local: production registrations are fixed, never plugins. - adapters[kind] = Adapter{Mode: "nodes", ValidateSpecification: func(sandbox.DeploymentSpec) error { return nil }, Normalize: nodeSelection(func(sandbox.DeploymentSpec) error { return nil })} + adapters[kind] = Adapter{Mode: "nodes", Operations: docker.Operations, ValidateSpecification: func(sandbox.DeploymentSpec) error { return nil }, Normalize: nodeSelection(func(sandbox.DeploymentSpec) error { return nil })} defer delete(adapters, kind) s, err := Normalize(sandbox.Selection{Provider: kind}) if err != nil || s.Provider != kind || !IsNode(kind) || SupportsCheckpoint(kind) { @@ -82,7 +84,11 @@ func TestRetainedLimitUsesRegisteredCapabilities(t *testing.T) { const kind = "capacity-test-provider" defer delete(adapters, kind) for _, checkpoint := range []bool{false, true} { - adapters[kind] = Adapter{Mode: "nodes", Checkpoint: checkpoint} + operations := docker.Operations + if checkpoint { + operations = microsandbox.Operations + } + adapters[kind] = Adapter{Mode: "nodes", Operations: operations} for _, retained := range []int{0, 20} { want := 10 if checkpoint { diff --git a/services/agents-api/internal/sandbox/sandbox_provider.go b/services/agents-api/internal/sandbox/sandbox_provider.go index 531fd3d71..170e402bb 100644 --- a/services/agents-api/internal/sandbox/sandbox_provider.go +++ b/services/agents-api/internal/sandbox/sandbox_provider.go @@ -4,10 +4,11 @@ // SandboxProvider owns compute and bootstrap, Runtime owns capability preparation, // and Harness adapters own native execution. Compute running is not execution ready. // -// Required operations are on SandboxProvider. CheckpointProvider is an independent -// optional capability; runtimeobs.Source and runtimeobs.BatchSource define optional -// read-only observation. Implementing an interface is a capability claim, not proof -// of qualification: run the common contract tests and adapter-specific acceptance. +// Required operations are on SandboxProvider. CheckpointProvider and runtimeobs +// observation remain separate small interfaces. Every registered adapter explicitly +// declares and implements each operation, including safe Unsupported rejections. +// Method-set presence never means an extension is supported. ValidateProvider and +// the common contract tests check declaration completeness and implementation. // // Registration is explicit construction, not a global init-time registry. Node-local // adapters register in sandbox/providers; Core's managed setup @@ -20,6 +21,7 @@ package sandbox import ( "context" "errors" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" ) var ( @@ -76,6 +78,7 @@ type CommandResult struct { // Implementations verify installation plus Reference ownership before mutation. // See docs/sandbox-provider.md for settlement, cleanup and retry requirements. type SandboxProvider interface { + providercontract.Declared // Create performs the original attempt once; no credential overwrite on conflict. Create(context.Context, Bootstrap) (Info, error) // GetInfo observes compute without creating, starting, renewing or preparing it. @@ -89,7 +92,7 @@ type SandboxProvider interface { RunCommand(context.Context, Reference, Command) (CommandResult, error) } -// CheckpointProvider is optional. It supplements the existing provider with +// CheckpointProvider is an explicitly declared extension. It supplies // exact-incarnation operations; Worker and Store remain the lifecycle owner. type CheckpointProvider interface { SandboxProvider diff --git a/services/agents-api/internal/store/provider_operations_fixture_test.go b/services/agents-api/internal/store/provider_operations_fixture_test.go new file mode 100644 index 000000000..b285590b3 --- /dev/null +++ b/services/agents-api/internal/store/provider_operations_fixture_test.go @@ -0,0 +1,92 @@ +package store_test + +import ( + "context" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeobs" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" +) + +func (*lifecycleProvider) ProviderOperations() providercontract.Operations { + return providercontract.Operations{ + "Create": {State: providercontract.Supported}, + "GetInfo": {State: providercontract.Supported}, + "Renew": {State: providercontract.Supported}, + "Kill": {State: providercontract.Supported}, + "RunCommand": {State: providercontract.Supported}, + "Initial": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "NewCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "GetCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "Suspend": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "Resume": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "KillCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "DeleteSnapshot": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "RunCommandCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "ResumeCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "Observe": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "ObserveBatch": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "DiscoverSelection": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "VerifyCredential": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + } +} +func (*lifecycleProvider) Initial(context.Context, sandbox.Reference) (sandbox.Compute, error) { + return sandbox.Compute{}, &providercontract.UnsupportedError{Operation: "Initial", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleProvider) NewCompute(context.Context, sandbox.Reference, uint64, *sandbox.SnapshotIdentity) (sandbox.Compute, error) { + return sandbox.Compute{}, &providercontract.UnsupportedError{Operation: "NewCompute", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleProvider) GetCompute(context.Context, sandbox.Reference, sandbox.Compute) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "GetCompute", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleProvider) Suspend(context.Context, sandbox.SuspendRequest) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "Suspend", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleProvider) Resume(context.Context, sandbox.ResumeRequest) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "Resume", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleProvider) KillCompute(context.Context, sandbox.Reference, sandbox.Compute) error { + return &providercontract.UnsupportedError{Operation: "KillCompute", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleProvider) DeleteSnapshot(context.Context, sandbox.Reference, sandbox.SnapshotIdentity) error { + return &providercontract.UnsupportedError{Operation: "DeleteSnapshot", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleProvider) RunCommandCompute(context.Context, sandbox.Reference, sandbox.Compute, sandbox.Command) (sandbox.CommandResult, error) { + return sandbox.CommandResult{}, &providercontract.UnsupportedError{Operation: "RunCommandCompute", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleProvider) ResumeCompute(context.Context, sandbox.Reference, sandbox.Compute) (sandbox.ComputeState, error) { + return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "ResumeCompute", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleProvider) Observe(context.Context, runtimeobs.Target) (runtimeobs.Sample, error) { + return runtimeobs.Sample{}, &providercontract.UnsupportedError{Operation: "Observe", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleProvider) ObserveBatch(context.Context, []runtimeobs.Target) ([]runtimeobs.BatchResult, error) { + return nil, &providercontract.UnsupportedError{Operation: "ObserveBatch", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleProvider) DiscoverSelection(context.Context, sandbox.Selection) (sandbox.Selection, error) { + return sandbox.Selection{}, &providercontract.UnsupportedError{Operation: "DiscoverSelection", Reason: "fixture_operation_not_supported"} +} +func (*lifecycleProvider) VerifyCredential(context.Context, []sandbox.Reference) error { + return &providercontract.UnsupportedError{Operation: "VerifyCredential", Reason: "fixture_operation_not_supported"} +} +func (*fakeCheckpointProvider) ProviderOperations() providercontract.Operations { + return providercontract.Operations{ + "Create": {State: providercontract.Supported}, + "GetInfo": {State: providercontract.Supported}, + "Renew": {State: providercontract.Supported}, + "Kill": {State: providercontract.Supported}, + "RunCommand": {State: providercontract.Supported}, + "Initial": {State: providercontract.Supported}, + "NewCompute": {State: providercontract.Supported}, + "GetCompute": {State: providercontract.Supported}, + "Suspend": {State: providercontract.Supported}, + "Resume": {State: providercontract.Supported}, + "KillCompute": {State: providercontract.Supported}, + "DeleteSnapshot": {State: providercontract.Supported}, + "RunCommandCompute": {State: providercontract.Supported}, + "ResumeCompute": {State: providercontract.Supported}, + "Observe": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "ObserveBatch": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "DiscoverSelection": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + "VerifyCredential": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"}, + } +} From c7627bfb6cf0f1aaca9ad9ea38557419f3310ecf Mon Sep 17 00:00:00 2001 From: saladday <1203511142@qq.com> Date: Wed, 30 Sep 2026 05:48:00 +0000 Subject: [PATCH 2/2] fix: retain selection discovery contract errors --- services/agents-api/cmd/server/managed_setup.go | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/services/agents-api/cmd/server/managed_setup.go b/services/agents-api/cmd/server/managed_setup.go index 2b632f911..1167b664d 100644 --- a/services/agents-api/cmd/server/managed_setup.go +++ b/services/agents-api/cmd/server/managed_setup.go @@ -4,12 +4,12 @@ import ( "context" "errors" "fmt" - "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" "sync/atomic" "time" "github.com/MiniMax-AI-Dev/parsar/internal/obs/log" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/execution" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/providercontract" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeobs" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox/node" @@ -106,6 +106,8 @@ func (s *managedSetup) prepare(ctx context.Context, setup store.SandboxSetup) (e if err != nil { return execution.PreparedRuntimeDeployment{}, err } + } else if !errors.Is(err, providercontract.ErrUnsupported) { + return execution.PreparedRuntimeDeployment{}, err } candidate.Selection = &selection return s.routeGenerations(candidate, setup)