diff --git a/contracts/agents-api/runtime-observability.md b/contracts/agents-api/runtime-observability.md index 078bccd2d..9394c6472 100644 --- a/contracts/agents-api/runtime-observability.md +++ b/contracts/agents-api/runtime-observability.md @@ -20,6 +20,10 @@ Managed Docker, microsandbox and E2B allocations are observed. `none` and `self_ The allocation's persisted `provider_key` selects exactly one configured source, which verifies the allocation's labels or equivalent ownership data before it returns values. Before any provider read, the allocation state decides some rows: `creating` or no allocation yet gives `allocation_pending`, `cleanup_pending` or `released` gives `runtime_not_running`, and a provider key without a source gives `source_not_configured`. A provider read that exceeds its deadline gives `sample_timeout`, a not-running result `runtime_not_running`, and an unavailable result `sample_unavailable`. Any other error, an ownership mismatch or an invalid sample fails the read. +The observation boundary is declared in `services/core/internal/runtimeobs/source.go`. Each registered `SourceResolver` declares supported `ResolveObservationSource`; registration validates this declaration without loading configuration or reading the database. Core resolves each provider key once per page, then validates the returned `Source` and uses that same immutable source for every read of the key on that page. An unconfigured resolver returns typed `ErrUnavailable`, which produces `sample_unavailable` without a provider type. Other resolution errors follow the provider-read error rules above. + +A source declares supported `ObservationProviderType`, which returns its immutable telemetry identity: a lowercase letter followed by at most 31 lowercase letters, digits or underscores. Empty identities are invalid. Provider registration and every resolved binding validate the identity and operation declarations before any sample or export. Reconfiguration affects later source resolutions; it cannot change the identity or provider selected for an in-flight page. Generation routers continue to resolve each allocation through its recorded deployment generation. + A source implements `Observe` and declares `ObserveBatch` in its provider operations. When `ObserveBatch` is declared supported, one call reads up to 100 targets of that provider; when it is declared unsupported, Core reads each target with `Observe`. A failed batch read is never retried target by target. The [Sandbox Provider guide](../../docs/sandbox-provider.md) describes the operation declarations. ## Sample semantics diff --git a/services/core/cmd/server/main.go b/services/core/cmd/server/main.go index 202fb1258..3f65cac32 100644 --- a/services/core/cmd/server/main.go +++ b/services/core/cmd/server/main.go @@ -38,6 +38,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/databaseurl" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/nativeinstaller" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/providercontract" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtime" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeenrollment" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimehistory" @@ -145,13 +146,15 @@ func run() error { if managedNodes != nil { managed = managedNodes.runtime } - observationSources := map[string]runtimeobs.Source{} + observationSources := map[string]runtimeobs.SourceResolver{} if managedNodes != nil && managedNodes.setup != nil { observationSources[managed.InstallationID] = managedNodes.setup } else if managed != nil { - if source, ok := managed.Provider.(runtimeobs.Source); ok { - observationSources[managed.InstallationID] = source + source, ok := managed.Provider.(runtimeobs.SourceResolver) + if !ok { + return providercontract.ErrContract } + observationSources[managed.InstallationID] = source } observationResolver, err := observationstoreresolver.NewResolver(executionStore) if err != nil { diff --git a/services/core/cmd/server/managed_generations.go b/services/core/cmd/server/managed_generations.go index 27bce505f..e77a5520d 100644 --- a/services/core/cmd/server/managed_generations.go +++ b/services/core/cmd/server/managed_generations.go @@ -26,9 +26,10 @@ 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 - operations providercontract.Operations + setup *managedSetup + store generationStore + operations providercontract.Operations + providerType string } func (p *generationRouter) route(ctx context.Context, r sandbox.Reference) (sandbox.SandboxProvider, func(), error) { @@ -91,6 +92,11 @@ func (p *generationRouter) RunCommand(ctx context.Context, r sandbox.Reference, type observedGenerationRouter struct{ *generationRouter } +func (p *observedGenerationRouter) ObservationProviderType() string { return p.providerType } +func (p *observedGenerationRouter) ResolveObservationSource(context.Context) (runtimeobs.Source, error) { + return p, nil +} + func (p *generationRouter) ProviderOperations() providercontract.Operations { return maps.Clone(p.operations) } @@ -108,6 +114,9 @@ func (p *observedGenerationRouter) Observe(ctx context.Context, t runtimeobs.Tar if !ok { return runtimeobs.Sample{}, providercontract.ErrContract } + if source.ObservationProviderType() != p.providerType { + return runtimeobs.Sample{}, providercontract.ErrContract + } return source.Observe(ctx, t) } @@ -127,7 +136,7 @@ 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, operations: candidate.Config.Provider.ProviderOperations()} + router := &generationRouter{setup: s, store: db, providerType: candidate.Config.Provider.(runtimeobs.Source).ObservationProviderType(), 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"} diff --git a/services/core/cmd/server/managed_setup.go b/services/core/cmd/server/managed_setup.go index cd6f55055..3ec771c93 100644 --- a/services/core/cmd/server/managed_setup.go +++ b/services/core/cmd/server/managed_setup.go @@ -134,18 +134,10 @@ func (s *managedSetup) configuration(setup store.SandboxSetup) (execution.Prepar } 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 - } - return "" + return providercontract.Operations{"ResolveObservationSource": {State: providercontract.Supported}} } -// 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) { +func (s *managedSetup) ResolveObservationSource(ctx context.Context) (runtimeobs.Source, error) { selected, err := s.load(ctx) if err != nil { return nil, err @@ -153,32 +145,11 @@ func (s *managedSetup) ObserveBatch(ctx context.Context, targets []runtimeobs.Ta 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, providercontract.ErrContract - } - return source.ObserveBatch(ctx, targets) -} - -func (s *managedSetup) Observe(ctx context.Context, target runtimeobs.Target) (runtimeobs.Sample, error) { - selected, err := s.load(ctx) - if err != nil { - return runtimeobs.Sample{}, err - } - 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{}, providercontract.ErrContract + return nil, providercontract.ErrContract } - return source.Observe(ctx, target) + return source, nil } func (s *managedSetup) provider(setup store.SandboxSetup) (sandbox.SandboxProvider, error) { diff --git a/services/core/cmd/server/managed_setup_test.go b/services/core/cmd/server/managed_setup_test.go index ea6cef8ea..d43376144 100644 --- a/services/core/cmd/server/managed_setup_test.go +++ b/services/core/cmd/server/managed_setup_test.go @@ -10,7 +10,11 @@ import ( "strings" "testing" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/providercontract" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox/docker" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox/e2b" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox/microsandbox" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" @@ -172,7 +176,7 @@ func TestManagedSetupResetTombstoneRejectsDelayedProviderLoad(t *testing.T) { if err := <-done; err != nil { t.Fatal(err) } - if s.selected.Load().Generation != 2 || s.selected.Load().Config != nil || s.ObservationProviderType() != "" { + if s.selected.Load().Generation != 2 || s.selected.Load().Config != nil { t.Fatal("empty publication lost its generation") } s.publish(&execution.RuntimeProvider{Generation: 1, ProviderKind: "docker"}) @@ -181,7 +185,7 @@ func TestManagedSetupResetTombstoneRejectsDelayedProviderLoad(t *testing.T) { } next := &execution.RuntimeProvider{Generation: 3, ProviderKind: "microsandbox"} s.publish(next) - if s.selected.Load().Config != next || s.ObservationProviderType() != "microsandbox" { + if s.selected.Load().Config != next { t.Fatal("reset blocked subsequent configuration") } } @@ -219,3 +223,40 @@ func testProviderPaths(t *testing.T, helper, state string) sandbox.ProcessPaths } return paths } + +func TestManagedObservationSourceKeepsSelectionAcrossReconfiguration(t *testing.T) { + db := &setupStore{value: store.SandboxSetup{InstallationID: "installation", Generation: 1}} + setup := &managedSetup{store: db, installationID: "installation"} + if source, err := setup.ResolveObservationSource(t.Context()); source != nil || !errors.Is(err, runtimeobs.ErrUnavailable) { + t.Fatal("unconfigured setup did not return typed unavailability", source, err) + } + first := &docker.Provider{} + db.value.Provider, db.value.Generation = "docker", 2 + setup.publish(&execution.RuntimeProvider{Generation: 2, Provider: first}) + source, err := setup.ResolveObservationSource(t.Context()) + if err != nil || source != first { + t.Fatal(source, err) + } + next := µsandbox.Provider{} + db.value.Provider, db.value.Generation = "microsandbox", 3 + setup.publish(&execution.RuntimeProvider{Generation: 3, Provider: next}) + if source.ObservationProviderType() != "docker" { + t.Fatal("in-flight identity changed") + } + selected, err := setup.ResolveObservationSource(t.Context()) + if err != nil || selected != next || selected.ObservationProviderType() != "microsandbox" { + t.Fatal(selected, err) + } +} + +func TestObservationGenerationIdentityMustMatchRoutedAllocation(t *testing.T) { + hub := node.NewHub(node.HubOptions{}) + defer hub.Close() + db := &setupStore{value: store.SandboxSetup{Provider: "docker"}} + setup := &managedSetup{hub: hub, store: db} + source := &observedGenerationRouter{&generationRouter{setup: setup, store: db, providerType: "e2b"}} + _, err := source.Observe(t.Context(), runtimeobs.Target{TenantID: "tenant", EnvironmentID: "environment", Instance: runtimeobs.Instance{AllocationID: "allocation"}}) + if !errors.Is(err, providercontract.ErrContract) { + t.Fatal("routed allocation was attributed to another provider", err) + } +} diff --git a/services/core/internal/execution/provider_operations_fixture_test.go b/services/core/internal/execution/provider_operations_fixture_test.go index 17050ff85..4b568704b 100644 --- a/services/core/internal/execution/provider_operations_fixture_test.go +++ b/services/core/internal/execution/provider_operations_fixture_test.go @@ -9,24 +9,26 @@ import ( 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"}, + "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"}, + "ObservationProviderType": {State: providercontract.Supported}, + "ResolveObservationSource": {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"}, } } func (*lifecycleOnlySandbox) Initial(context.Context, sandbox.Reference) (sandbox.Compute, error) { @@ -68,3 +70,8 @@ func (*lifecycleOnlySandbox) DiscoverSelection(context.Context, sandbox.Selectio func (*lifecycleOnlySandbox) VerifyCredential(context.Context, []sandbox.Reference) error { return &providercontract.UnsupportedError{Operation: "VerifyCredential", Reason: "fixture_operation_not_supported"} } + +func (*lifecycleOnlySandbox) ObservationProviderType() string { return "fixture" } +func (p *lifecycleOnlySandbox) ResolveObservationSource(context.Context) (runtimeobs.Source, error) { + return p, nil +} diff --git a/services/core/internal/runtimeobs/exporter_independence_test.go b/services/core/internal/runtimeobs/exporter_independence_test.go index add9f2276..2c3b92078 100644 --- a/services/core/internal/runtimeobs/exporter_independence_test.go +++ b/services/core/internal/runtimeobs/exporter_independence_test.go @@ -10,7 +10,7 @@ func TestExporterOutageDoesNotDelayOtherDestinations(t *testing.T) { target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} blocked := &gatedExporter{started: make(chan struct{}, 1), release: make(chan struct{})} records := make(chan ExportRecord, 1) - service, err := NewService(fixedResolver{target: target}, map[string]Source{"provider": &fixedSource{sample: Sample{ObservedAt: now}}}, + service, err := NewService(fixedResolver{target: target}, map[string]SourceResolver{"provider": &fixedSource{sample: Sample{ObservedAt: now}}}, WithExporter(blocked, ExportOptions{QueueCapacity: 1, Timeout: time.Second}), WithExporter(channelExporter{records: records}, ExportOptions{QueueCapacity: 1, Timeout: time.Second})) if err != nil { diff --git a/services/core/internal/runtimeobs/operations_fixture_test.go b/services/core/internal/runtimeobs/operations_fixture_test.go index 2bb1ccd06..d55ebe8c4 100644 --- a/services/core/internal/runtimeobs/operations_fixture_test.go +++ b/services/core/internal/runtimeobs/operations_fixture_test.go @@ -6,23 +6,42 @@ import ( ) func (*fixedSource) ProviderOperations() providercontract.Operations { - return providercontract.Operations{"Observe": {State: providercontract.Supported}, "ObserveBatch": {State: providercontract.Unsupported, Reason: "fixture_has_no_batch_observation"}} + return providercontract.Operations{"ResolveObservationSource": {State: providercontract.Supported}, "ObservationProviderType": {State: providercontract.Supported}, "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"}} + return providercontract.Operations{"ResolveObservationSource": {State: providercontract.Supported}, "ObservationProviderType": {State: providercontract.Supported}, "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"}} + return providercontract.Operations{"ResolveObservationSource": {State: providercontract.Supported}, "ObservationProviderType": {State: providercontract.Supported}, "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}} + return providercontract.Operations{"ResolveObservationSource": {State: providercontract.Supported}, "ObservationProviderType": {State: providercontract.Supported}, "Observe": {State: providercontract.Supported}, "ObserveBatch": {State: providercontract.Supported}} +} + +func (s *fixedSource) ResolveObservationSource(context.Context) (Source, error) { return s, nil } +func (*fixedSource) ObservationProviderType() string { return "fixture" } + +func (s blockingSource) ResolveObservationSource(context.Context) (Source, error) { return s, nil } + +func (s *countingSource) ResolveObservationSource(context.Context) (Source, error) { return s, nil } +func (*countingSource) ObservationProviderType() string { return "fixture" } + +func (s *batchSource) ResolveObservationSource(context.Context) (Source, error) { return s, nil } +func (*batchSource) ObservationProviderType() string { return "fixture" } + +func (s typedSource) ResolveObservationSource(context.Context) (Source, error) { return s, nil } + +func (s *failingBatchSource) ResolveObservationSource(context.Context) (Source, error) { return s, nil } + +func (s *unsupportedObservation) ResolveObservationSource(context.Context) (Source, error) { + return s, nil } diff --git a/services/core/internal/runtimeobs/operations_test.go b/services/core/internal/runtimeobs/operations_test.go index 179c6e6b3..3e66de33c 100644 --- a/services/core/internal/runtimeobs/operations_test.go +++ b/services/core/internal/runtimeobs/operations_test.go @@ -15,7 +15,7 @@ type failingBatchSource struct { } func (*failingBatchSource) ProviderOperations() providercontract.Operations { - return providercontract.Operations{"Observe": {State: providercontract.Supported}, "ObserveBatch": {State: providercontract.Supported}} + return providercontract.Operations{"ResolveObservationSource": {State: providercontract.Supported}, "ObservationProviderType": {State: providercontract.Supported}, "Observe": {State: providercontract.Supported}, "ObserveBatch": {State: providercontract.Supported}} } func (s *failingBatchSource) ObserveBatch(context.Context, []Target) ([]BatchResult, error) { s.batches++ @@ -40,7 +40,7 @@ func TestBatchFallbackRequiresExplicitSafeUnsupported(t *testing.T) { 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}) + service, err := NewService(fixedResolver{target: target}, map[string]SourceResolver{"provider": source}) if err != nil { t.Fatal(err) } @@ -58,12 +58,12 @@ func TestBatchFallbackRequiresExplicitSafeUnsupported(t *testing.T) { 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"}} + return providercontract.Operations{"ResolveObservationSource": {State: providercontract.Supported}, "ObservationProviderType": {State: providercontract.Supported}, "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}) + service, err := NewService(fixedResolver{target: target}, map[string]SourceResolver{"provider": source}) if err != nil { t.Fatal(err) } diff --git a/services/core/internal/runtimeobs/sampler_test.go b/services/core/internal/runtimeobs/sampler_test.go index 80a96de30..9b62231eb 100644 --- a/services/core/internal/runtimeobs/sampler_test.go +++ b/services/core/internal/runtimeobs/sampler_test.go @@ -146,7 +146,7 @@ func TestSamplerSweepsEveryPageAndIsolatesSessionFailures(t *testing.T) { func TestSamplerBoundsConcurrencyAndSourceDeadline(t *testing.T) { source := &countingSource{} 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}) + service, err := NewService(fixedResolver{target: target}, map[string]SourceResolver{"provider": source}) if err != nil { t.Fatal(err) } @@ -302,7 +302,7 @@ func TestSamplerPreservesProviderTimeoutAndFinalFenceAfterSlowResolution(t *test records := make(chan ExportRecord, 1) service, err := NewService( resolver, - map[string]Source{"provider": blockingSource{}}, + map[string]SourceResolver{"provider": blockingSource{}}, WithExporter(channelExporter{records: records}, ExportOptions{}), ) if err != nil { diff --git a/services/core/internal/runtimeobs/service.go b/services/core/internal/runtimeobs/service.go index a41f363d7..101c6fe9f 100644 --- a/services/core/internal/runtimeobs/service.go +++ b/services/core/internal/runtimeobs/service.go @@ -4,14 +4,12 @@ import ( "context" "errors" "fmt" - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/providercontract" "reflect" - "regexp" "sync" "time" -) -var providerTypePattern = regexp.MustCompile(`^[a-z][a-z0-9_]{0,31}$`) + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/providercontract" +) type Observation struct { Target Target @@ -32,12 +30,12 @@ const ( type Service struct { resolver TargetResolver - sources map[string]Source + sources map[string]SourceResolver now func() time.Time exports []*exportDispatcher } -func NewService(resolver TargetResolver, sources map[string]Source, options ...ServiceOption) (*Service, error) { +func NewService(resolver TargetResolver, sources map[string]SourceResolver, options ...ServiceOption) (*Service, error) { if resolver == nil { return nil, errors.New("Runtime observation resolver is required") } @@ -50,14 +48,17 @@ func NewService(resolver TargetResolver, sources map[string]Source, options ...S return nil, err } } - copySources := make(map[string]Source, len(sources)) + copySources := make(map[string]SourceResolver, len(sources)) for key, source := range sources { 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 { + if err := providercontract.Validate(source, reflect.TypeFor[SourceResolver]()); err != nil { return nil, err } + if err := providercontract.Require(source, "ResolveObservationSource"); err != nil { + return nil, fmt.Errorf("%w: observation source resolution must be supported", providercontract.ErrContract) + } copySources[key] = source } service := &Service{resolver: resolver, sources: copySources, now: time.Now} @@ -175,6 +176,22 @@ func (s *Service) observeSessions(ctx context.Context, sessions []SessionIdentit } for _, key := range keys { group := groups[key] + sourceCtx, stop := sourceContext(ctx, options.SourceTimeout) + source, err := s.sources[key].ResolveObservationSource(sourceCtx) + stop() + if err == nil { + err = ValidateSource(source) + } + if err != nil { + for _, read := range group { + observations[read.index], errs[read.index] = s.complete(ctx, read, Sample{}, err, 0, collectionSource, owner) + } + continue + } + providerType := source.ObservationProviderType() + for _, read := range group { + read.source, read.providerType = source, providerType + } for start := 0; start < len(group); start += MaxBatchTargets { chunk := group[start:min(start+MaxBatchTargets, len(group))] if s.readBatch(ctx, chunk, observations, errs, collectionSource, owner, options.SourceTimeout) { @@ -300,18 +317,11 @@ func (s *Service) resolve(ctx context.Context, tenantID, sessionID string, owner default: return Observation{}, nil, errors.New("invalid managed Runtime allocation state") } - source, ok := s.sources[target.Instance.ProviderKey] + _, ok := s.sources[target.Instance.ProviderKey] if !ok { return Observation{Target: target, Status: StatusUnavailable, Reason: "source_not_configured", ResolvedAt: resolvedAt}, nil, nil } - providerType := "" - if typed, ok := source.(interface{ ObservationProviderType() string }); ok { - providerType = typed.ObservationProviderType() - if providerType != "" && !providerTypePattern.MatchString(providerType) { - return Observation{}, nil, errors.New("invalid Runtime observation provider type") - } - } - return Observation{}, &sourceRead{key: target.Instance.ProviderKey, source: source, target: target, providerType: providerType}, nil + return Observation{}, &sourceRead{key: target.Instance.ProviderKey, target: target}, nil } // complete classifies one provider result and hands it to history export. diff --git a/services/core/internal/runtimeobs/service_test.go b/services/core/internal/runtimeobs/service_test.go index 51d5b9e59..ca3d9dbc2 100644 --- a/services/core/internal/runtimeobs/service_test.go +++ b/services/core/internal/runtimeobs/service_test.go @@ -125,7 +125,7 @@ func TestServiceDoesNotCallSourcesForUnsupportedModes(t *testing.T) { if mode == ModeSelfHosted { target.EnvironmentID = "environment" } - service, err := NewService(fixedResolver{target: target}, map[string]Source{"provider": source}) + service, err := NewService(fixedResolver{target: target}, map[string]SourceResolver{"provider": source}) if err != nil { t.Fatal(err) } @@ -151,7 +151,7 @@ func TestServicePreservesUnavailableAndObservedZero(t *testing.T) { zeroMemory := uint64(0) now := time.Date(2026, 9, 22, 1, 0, 0, 0, time.UTC) source := &fixedSource{sample: Sample{ObservedAt: now, CPUUsageSecondsTotal: &zeroCPU, MemoryUsageBytes: &zeroMemory}} - service, err = NewService(fixedResolver{target: target}, map[string]Source{"provider": source}) + service, err = NewService(fixedResolver{target: target}, map[string]SourceResolver{"provider": source}) if err != nil { t.Fatal(err) } @@ -182,7 +182,7 @@ func TestServiceExportsOnlySanitizedValidatedRecords(t *testing.T) { MemoryUsageBytes: &memoryUsage, MemoryLimitBytes: &memoryLimit, }}, providerType: "docker"} records := make(chan ExportRecord, 1) - service, err := NewService(fixedResolver{target: target}, map[string]Source{"provider": source}, WithExporter(channelExporter{records: records}, ExportOptions{QueueCapacity: 1, Timeout: time.Second})) + service, err := NewService(fixedResolver{target: target}, map[string]SourceResolver{"provider": source}, WithExporter(channelExporter{records: records}, ExportOptions{QueueCapacity: 1, Timeout: time.Second})) if err != nil { t.Fatal(err) } @@ -224,7 +224,7 @@ func TestServiceMarksPeriodicHistoryCollection(t *testing.T) { records := make(chan ExportRecord, 1) service, err := NewService( fixedResolver{target: target}, - map[string]Source{"provider": &fixedSource{sample: Sample{ObservedAt: now, StartedAt: &startedAt}}}, + map[string]SourceResolver{"provider": &fixedSource{sample: Sample{ObservedAt: now, StartedAt: &startedAt}}}, WithExporter(channelExporter{records: records}, ExportOptions{}), ) if err != nil { @@ -257,7 +257,7 @@ func TestServiceDoesNotExportPeriodicSampleAfterOwnershipLoss(t *testing.T) { records := make(chan ExportRecord, 1) service, err := NewService( fixedResolver{target: target}, - map[string]Source{"provider": &fixedSource{sample: Sample{ObservedAt: now, StartedAt: &startedAt}}}, + map[string]SourceResolver{"provider": &fixedSource{sample: Sample{ObservedAt: now, StartedAt: &startedAt}}}, WithExporter(channelExporter{records: records}, ExportOptions{}), ) if err != nil { @@ -285,7 +285,7 @@ func TestServiceRejectsUnsafeProviderTypeBeforeSamplingOrExport(t *testing.T) { records := make(chan ExportRecord, 1) service, err := NewService( fixedResolver{target: target}, - map[string]Source{"provider": source}, + map[string]SourceResolver{"provider": source}, WithExporter(channelExporter{records: records}, ExportOptions{}), ) if err != nil { @@ -311,7 +311,7 @@ func TestServiceExportQueueNeverBlocksOrChangesObservation(t *testing.T) { exporter := &gatedExporter{started: make(chan struct{}, 1), release: make(chan struct{})} service, err := NewService( fixedResolver{target: target}, - map[string]Source{"provider": &fixedSource{sample: Sample{ObservedAt: now}}}, + map[string]SourceResolver{"provider": &fixedSource{sample: Sample{ObservedAt: now}}}, WithExporter(exporter, ExportOptions{QueueCapacity: 1, Timeout: time.Second}), ) if err != nil { @@ -370,7 +370,7 @@ func TestServiceIgnoresExporterFailureAndValidatesOptions(t *testing.T) { target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} service, err := NewService( fixedResolver{target: target}, - map[string]Source{"provider": &fixedSource{sample: Sample{ObservedAt: now}}}, + map[string]SourceResolver{"provider": &fixedSource{sample: Sample{ObservedAt: now}}}, WithExporter(channelExporter{records: records, err: errors.New("backend unavailable")}, ExportOptions{}), ) if err != nil { @@ -392,7 +392,7 @@ func TestServiceCloseHonorsItsDeadlineWhenExporterDoesNot(t *testing.T) { target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} service, err := NewService( fixedResolver{target: target}, - map[string]Source{"provider": &fixedSource{sample: Sample{ObservedAt: now}}}, + map[string]SourceResolver{"provider": &fixedSource{sample: Sample{ObservedAt: now}}}, WithExporter(exporter, ExportOptions{QueueCapacity: 1, Timeout: time.Millisecond}), ) if err != nil { @@ -422,7 +422,7 @@ func TestServiceIsolatesExporterPanics(t *testing.T) { target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} service, err := NewService( fixedResolver{target: target}, - map[string]Source{"provider": &fixedSource{sample: Sample{ObservedAt: now}}}, + map[string]SourceResolver{"provider": &fixedSource{sample: Sample{ObservedAt: now}}}, WithExporter(exporter, ExportOptions{}), ) if err != nil { @@ -455,7 +455,7 @@ func TestServiceMapsOnlyDeclaredUnavailability(t *testing.T) { {err: context.DeadlineExceeded, wantReason: "sample_timeout"}, {err: errors.New("Docker permission denied"), wantError: true}, } { - service, err := NewService(fixedResolver{target: target}, map[string]Source{ + service, err := NewService(fixedResolver{target: target}, map[string]SourceResolver{ "provider": typedSource{fixedSource: &fixedSource{err: tc.err}, providerType: "docker"}, }) if err != nil { @@ -476,7 +476,7 @@ func TestServiceMapsAnActualSourceDeadlineWithoutLeakingIt(t *testing.T) { EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}, } - service, err := NewService(fixedResolver{target: target}, map[string]Source{"provider": blockingSource{}}) + service, err := NewService(fixedResolver{target: target}, map[string]SourceResolver{"provider": blockingSource{}}) if err != nil { t.Fatal(err) } @@ -496,7 +496,7 @@ func TestServiceExportsPeriodicSourceTimeoutAfterFinalOwnershipFence(t *testing. records := make(chan ExportRecord, 1) service, err := NewService( fixedResolver{target: target}, - map[string]Source{"provider": blockingSource{}}, + map[string]SourceResolver{"provider": blockingSource{}}, WithExporter(channelExporter{records: records}, ExportOptions{}), ) if err != nil { @@ -537,7 +537,7 @@ func TestServiceClassifiesResolverAndTerminalAllocationUnavailability(t *testing source := &fixedSource{} target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "creating"}} - service, err = NewService(fixedResolver{target: target}, map[string]Source{"provider": source}) + service, err = NewService(fixedResolver{target: target}, map[string]SourceResolver{"provider": source}) if err != nil { t.Fatal(err) } @@ -547,7 +547,7 @@ func TestServiceClassifiesResolverAndTerminalAllocationUnavailability(t *testing } target = Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "released"}} - service, err = NewService(fixedResolver{target: target}, map[string]Source{"provider": &fixedSource{}}) + service, err = NewService(fixedResolver{target: target}, map[string]SourceResolver{"provider": &fixedSource{}}) if err != nil { t.Fatal(err) } @@ -560,7 +560,7 @@ func TestServiceClassifiesResolverAndTerminalAllocationUnavailability(t *testing target = Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{ AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running", AllocationCreatedAt: time.Now().Add(time.Hour), }} - service, err = NewService(fixedResolver{target: target}, map[string]Source{"provider": source}) + service, err = NewService(fixedResolver{target: target}, map[string]SourceResolver{"provider": source}) if err != nil { t.Fatal(err) } @@ -584,7 +584,7 @@ func TestServiceRejectsUnsafeProviderSamples(t *testing.T) { {ObservedAt: now, MemoryUsageBytes: &tooLarge}, {ObservedAt: now, MemoryLimitBytes: &tooLarge}, } { - service, err := NewService(fixedResolver{target: target}, map[string]Source{"provider": &fixedSource{sample: sample}}) + service, err := NewService(fixedResolver{target: target}, map[string]SourceResolver{"provider": &fixedSource{sample: sample}}) if err != nil { t.Fatal(err) } @@ -644,7 +644,7 @@ func TestServiceBatchesPageReadsWithinProviderLimit(t *testing.T) { ratio := 0.25 source := &batchSource{sample: Sample{ObservedAt: now, CPUUtilizationRatio: &ratio}} 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}) + service, err := NewService(fixedResolver{target: target}, map[string]SourceResolver{"provider": source}) if err != nil { t.Fatal(err) } diff --git a/services/core/internal/runtimeobs/source.go b/services/core/internal/runtimeobs/source.go index 3a31edae9..8886a42b1 100644 --- a/services/core/internal/runtimeobs/source.go +++ b/services/core/internal/runtimeobs/source.go @@ -3,6 +3,10 @@ package runtimeobs import ( "context" "errors" + "fmt" + "reflect" + "regexp" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/providercontract" ) @@ -11,10 +15,20 @@ var ( ErrNotRunning = errors.New("Runtime is not running") ) +// SourceResolver selects one immutable source for a page of provider reads. +// An unconfigured resolver returns ErrUnavailable. Registration never resolves +// a source; callers validate each selected source before reading it. +type SourceResolver interface { + providercontract.Declared + ResolveObservationSource(context.Context) (Source, error) +} + // 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 + // ObservationProviderType is a nonempty, immutable telemetry identity. + ObservationProviderType() string Observe(context.Context, Target) (Sample, error) } @@ -37,3 +51,19 @@ type BatchResult struct { type TargetResolver interface { Resolve(context.Context, string, string) (Target, error) } + +var providerTypePattern = regexp.MustCompile(`^[a-z][a-z0-9_]{0,31}$`) + +// ValidateSource checks every read operation and its immutable provider identity. +func ValidateSource(source Source) error { + if err := providercontract.Validate(source, reflect.TypeFor[Source](), reflect.TypeFor[BatchSource]()); err != nil { + return err + } + if err := providercontract.Require(source, "ObservationProviderType"); err != nil { + return fmt.Errorf("%w: observation identity must be supported", providercontract.ErrContract) + } + if !providerTypePattern.MatchString(source.ObservationProviderType()) { + return fmt.Errorf("%w: invalid Runtime observation provider type", providercontract.ErrContract) + } + return nil +} diff --git a/services/core/internal/runtimeobs/source_test.go b/services/core/internal/runtimeobs/source_test.go new file mode 100644 index 000000000..e65753320 --- /dev/null +++ b/services/core/internal/runtimeobs/source_test.go @@ -0,0 +1,137 @@ +package runtimeobs + +import ( + "context" + "errors" + "testing" + "time" + + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/providercontract" +) + +type selectingSource struct { + source Source + err error + calls int + support providercontract.Support +} + +func (s *selectingSource) ProviderOperations() providercontract.Operations { + return providercontract.Operations{"ResolveObservationSource": s.support} +} +func (s *selectingSource) ResolveObservationSource(context.Context) (Source, error) { + s.calls++ + return s.source, s.err +} + +type identityDeclaration struct { + typedSource + support providercontract.Support +} + +func (s identityDeclaration) ProviderOperations() providercontract.Operations { + operations := s.typedSource.ProviderOperations() + operations["ObservationProviderType"] = s.support + return operations +} + +func TestSourceBindingRejectsInvalidIdentityBeforeReadOrExport(t *testing.T) { + for _, test := range []struct { + name, identity string + support providercontract.Support + }{ + {"empty", "", providercontract.Support{State: providercontract.Supported}}, + {"unsafe", "secret:provider", providercontract.Support{State: providercontract.Supported}}, + {"undeclared", "docker", providercontract.Support{}}, + {"unsupported", "docker", providercontract.Support{State: providercontract.Unsupported, Reason: "missing_identity"}}, + } { + t.Run(test.name, func(t *testing.T) { + source := identityDeclaration{typedSource: typedSource{fixedSource: &fixedSource{}, providerType: test.identity}, support: test.support} + resolver := &selectingSource{source: source, support: providercontract.Support{State: providercontract.Supported}} + records := make(chan ExportRecord, 1) + service, err := NewService(observableTarget(), map[string]SourceResolver{"provider": resolver}, WithExporter(channelExporter{records: records}, ExportOptions{})) + if err != nil { + t.Fatal(err) + } + if _, err := service.ObserveSession(t.Context(), "tenant", "session"); !errors.Is(err, providercontract.ErrContract) { + t.Fatalf("invalid identity accepted: %v", err) + } + if source.calls != 0 { + t.Fatal("invalid identity reached provider read") + } + if err := service.Close(t.Context()); err != nil { + t.Fatal(err) + } + if len(records) != 0 { + t.Fatal("invalid identity reached export") + } + }) + } +} + +func observableTarget() fixedResolver { + return fixedResolver{target: Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}}} +} + +func TestSourceResolverRegistrationAndUnavailableSelection(t *testing.T) { + for _, support := range []providercontract.Support{{}, {State: providercontract.Unsupported, Reason: "no_source"}} { + resolver := &selectingSource{support: support} + if _, err := NewService(observableTarget(), map[string]SourceResolver{"provider": resolver}); !errors.Is(err, providercontract.ErrContract) || resolver.calls != 0 { + t.Fatal("invalid resolver registration accepted or performed I/O", err) + } + } + resolver := &selectingSource{err: ErrUnavailable, support: providercontract.Support{State: providercontract.Supported}} + service, err := NewService(observableTarget(), map[string]SourceResolver{"provider": resolver}) + if err != nil || resolver.calls != 0 { + t.Fatal("registration resolved unconfigured provider", err) + } + observation, err := service.ObserveSession(t.Context(), "tenant", "session") + if err != nil || observation.Status != StatusUnavailable || observation.Reason != "sample_unavailable" || observation.ProviderType != "" { + t.Fatal(observation, err) + } + var nilSource *fixedSource + resolver.err, resolver.source = nil, nilSource + if _, err := service.ObserveSession(t.Context(), "tenant", "session"); !errors.Is(err, providercontract.ErrContract) { + t.Fatal("typed nil source accepted", err) + } +} + +type reconfiguringSource struct { + typedSource + resolver *selectingSource + next Source +} + +func (s *reconfiguringSource) Observe(ctx context.Context, target Target) (Sample, error) { + s.resolver.source = s.next + return s.typedSource.Observe(ctx, target) +} + +func TestSourceSelectionStaysBoundForWholePage(t *testing.T) { + now := time.Now() + next := typedSource{fixedSource: &fixedSource{sample: Sample{ObservedAt: now}}, providerType: "next"} + resolver := &selectingSource{support: providercontract.Support{State: providercontract.Supported}} + previous := &reconfiguringSource{typedSource: typedSource{fixedSource: &fixedSource{sample: Sample{ObservedAt: now}}, providerType: "previous"}, resolver: resolver, next: next} + resolver.source = previous + service, err := NewService(observableTarget(), map[string]SourceResolver{"provider": resolver}) + if err != nil { + t.Fatal(err) + } + sessions := make([]SessionIdentity, MaxBatchTargets+1) + for index := range sessions { + sessions[index] = SessionIdentity{TenantID: "tenant", SessionID: "session"} + } + observations, errs := service.ObserveSessions(t.Context(), sessions, PageOptions{Concurrency: 1}) + for index, observation := range observations { + if errs[index] != nil || observation.ProviderType != "previous" || observation.Status != StatusObserved { + t.Fatal(index, observation, errs[index]) + } + } + if resolver.calls != 1 || previous.calls != len(sessions) || next.calls != 0 { + t.Fatalf("selection/read counts %d / %d / %d", resolver.calls, previous.calls, next.calls) + } + observation, err := service.ObserveSession(t.Context(), "tenant", "session") + if err != nil || observation.ProviderType != "next" || next.calls != 1 || resolver.calls != 2 { + t.Fatal("later page did not select new provider", observation, err) + } +} diff --git a/services/core/internal/sandbox/docker/operations.go b/services/core/internal/sandbox/docker/operations.go index 6135f7829..c2758008d 100644 --- a/services/core/internal/sandbox/docker/operations.go +++ b/services/core/internal/sandbox/docker/operations.go @@ -15,24 +15,26 @@ 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"}, + "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"}, + "ObservationProviderType": {State: providercontract.Supported}, + "ResolveObservationSource": {State: providercontract.Supported}, + "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() } diff --git a/services/core/internal/sandbox/docker/resources.go b/services/core/internal/sandbox/docker/resources.go index 387214dcb..d48dfcb17 100644 --- a/services/core/internal/sandbox/docker/resources.go +++ b/services/core/internal/sandbox/docker/resources.go @@ -91,3 +91,7 @@ func sampleFromDocker(inspected container.InspectResponse, stats dockerStatsResp } return sample, nil } + +func (p *Provider) ResolveObservationSource(context.Context) (runtimeobs.Source, error) { + return p, nil +} diff --git a/services/core/internal/sandbox/e2b/observations.go b/services/core/internal/sandbox/e2b/observations.go index 0b5da3c5f..e931805ad 100644 --- a/services/core/internal/sandbox/e2b/observations.go +++ b/services/core/internal/sandbox/e2b/observations.go @@ -160,3 +160,7 @@ func sampleFromObservation(observation Observation, now time.Time) (runtimeobs.S } func finite(value float64) bool { return !math.IsNaN(value) && !math.IsInf(value, 0) } + +func (p *Provider) ResolveObservationSource(context.Context) (runtimeobs.Source, error) { + return p, nil +} diff --git a/services/core/internal/sandbox/e2b/operations.go b/services/core/internal/sandbox/e2b/operations.go index f855e7d5f..14530ebf9 100644 --- a/services/core/internal/sandbox/e2b/operations.go +++ b/services/core/internal/sandbox/e2b/operations.go @@ -15,24 +15,26 @@ 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}, + "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"}, + "ObservationProviderType": {State: providercontract.Supported}, + "ResolveObservationSource": {State: providercontract.Supported}, + "Observe": {State: providercontract.Supported}, + "ObserveBatch": {State: providercontract.Supported}, + "DiscoverSelection": {State: providercontract.Supported}, + "VerifyCredential": {State: providercontract.Supported}, } } func (*Provider) ProviderOperations() providercontract.Operations { return Operations() } diff --git a/services/core/internal/sandbox/microsandbox/operations.go b/services/core/internal/sandbox/microsandbox/operations.go index 053ca9af2..f59c9ed02 100644 --- a/services/core/internal/sandbox/microsandbox/operations.go +++ b/services/core/internal/sandbox/microsandbox/operations.go @@ -15,24 +15,26 @@ 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"}, + "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}, + "ObservationProviderType": {State: providercontract.Supported}, + "ResolveObservationSource": {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() } diff --git a/services/core/internal/sandbox/microsandbox/resources.go b/services/core/internal/sandbox/microsandbox/resources.go index 77897ed65..6b4619147 100644 --- a/services/core/internal/sandbox/microsandbox/resources.go +++ b/services/core/internal/sandbox/microsandbox/resources.go @@ -74,3 +74,7 @@ func sampleFromMetrics(config Config, metrics Metrics) (runtimeobs.Sample, error MemoryUsageBytes: &memoryUsage, MemoryLimitBytes: &memoryLimit, }, nil } + +func (p *Provider) ResolveObservationSource(context.Context) (runtimeobs.Source, error) { + return p, nil +} diff --git a/services/core/internal/sandbox/node/observations.go b/services/core/internal/sandbox/node/observations.go index 253301087..515606640 100644 --- a/services/core/internal/sandbox/node/observations.go +++ b/services/core/internal/sandbox/node/observations.go @@ -50,3 +50,7 @@ func observeProvider(ctx context.Context, provider sandbox.SandboxProvider, targ } return source.Observe(ctx, target) } + +func (p *provider) ResolveObservationSource(context.Context) (runtimeobs.Source, error) { + return p, nil +} diff --git a/services/core/internal/sandbox/node/observations_test.go b/services/core/internal/sandbox/node/observations_test.go index bec9ec874..401ab9cf1 100644 --- a/services/core/internal/sandbox/node/observations_test.go +++ b/services/core/internal/sandbox/node/observations_test.go @@ -93,7 +93,7 @@ func TestObservationsRouteThroughAssignedNodeWithoutLifecycleCalls(t *testing.T) } return "", 0, sandbox.ErrOwnership }).(runtimeobs.Source) - if typed := source.(interface{ ObservationProviderType() string }).ObservationProviderType(); typed != "docker" { + if typed := source.ObservationProviderType(); typed != "docker" { t.Fatal("provider type lost", typed) } for _, test := range []struct { diff --git a/services/core/internal/sandbox/node/provider_operations_fixture_test.go b/services/core/internal/sandbox/node/provider_operations_fixture_test.go index 56c21119b..8d34c1511 100644 --- a/services/core/internal/sandbox/node/provider_operations_fixture_test.go +++ b/services/core/internal/sandbox/node/provider_operations_fixture_test.go @@ -9,24 +9,26 @@ import ( 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"}, + "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"}, + "ObservationProviderType": {State: providercontract.Supported}, + "ResolveObservationSource": {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"}, } } func (*fakeProvider) Initial(context.Context, sandbox.Reference) (sandbox.Compute, error) { @@ -70,23 +72,35 @@ func (*fakeProvider) VerifyCredential(context.Context, []sandbox.Reference) erro } 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"}, + "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"}, + "ObservationProviderType": {State: providercontract.Supported}, + "ResolveObservationSource": {State: providercontract.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"}, } } + +func (*fakeProvider) ObservationProviderType() string { return "fixture" } +func (p *fakeProvider) ResolveObservationSource(context.Context) (runtimeobs.Source, error) { + return p, nil +} + +func (*observationProvider) ObservationProviderType() string { return "fixture" } +func (p *observationProvider) ResolveObservationSource(context.Context) (runtimeobs.Source, error) { + return p, nil +} diff --git a/services/core/internal/sandbox/operations.go b/services/core/internal/sandbox/operations.go index 817a7baa7..3faf4454b 100644 --- a/services/core/internal/sandbox/operations.go +++ b/services/core/internal/sandbox/operations.go @@ -13,13 +13,16 @@ import ( var providerInterfaces = []reflect.Type{ reflect.TypeFor[SandboxProvider](), reflect.TypeFor[CheckpointProvider](), reflect.TypeFor[SelectionDiscoverer](), reflect.TypeFor[CredentialVerifier](), - reflect.TypeFor[runtimeobs.Source](), reflect.TypeFor[runtimeobs.BatchSource](), + reflect.TypeFor[runtimeobs.SourceResolver](), reflect.TypeFor[runtimeobs.Source](), reflect.TypeFor[runtimeobs.BatchSource](), } func ValidateProvider(p SandboxProvider) error { if err := providercontract.Validate(p, providerInterfaces...); err != nil { return err } + if err := runtimeobs.ValidateSource(p.(runtimeobs.Source)); err != nil { + return err + } return ValidateOperations(p.ProviderOperations()) } @@ -49,6 +52,11 @@ func ValidateOperations(operations providercontract.Operations) error { return fmt.Errorf("%w: required operation %s", providercontract.ErrContract, name) } } + for _, name := range []string{"ObservationProviderType", "ResolveObservationSource"} { + if 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]() diff --git a/services/core/internal/sandbox/operations_test.go b/services/core/internal/sandbox/operations_test.go index 0b0dafafa..f1cc091d1 100644 --- a/services/core/internal/sandbox/operations_test.go +++ b/services/core/internal/sandbox/operations_test.go @@ -28,6 +28,12 @@ func TestDeclarationsRejectMissingUnknownAndContradictoryOperations(t *testing.T name string change func(providercontract.Operations) }{ + {"identity unsupported", func(o providercontract.Operations) { + o["ObservationProviderType"] = providercontract.Support{State: providercontract.Unsupported, Reason: "no_identity"} + }}, + {"resolver unsupported", func(o providercontract.Operations) { + o["ResolveObservationSource"] = providercontract.Support{State: providercontract.Unsupported, Reason: "no_resolver"} + }}, {"omitted", func(o providercontract.Operations) { delete(o, "DeleteSnapshot") }}, {"zero", func(o providercontract.Operations) { o["ObserveBatch"] = providercontract.Support{} }}, {"unknown", func(o providercontract.Operations) { @@ -110,3 +116,12 @@ func TestNewContractRequiresAnAuthoredDecision(t *testing.T) { } } } + +type invalidObservationIdentity struct{ *docker.Provider } + +func (*invalidObservationIdentity) ObservationProviderType() string { return "" } +func TestProviderRegistrationRequiresObservationIdentity(t *testing.T) { + if err := sandbox.ValidateProvider(&invalidObservationIdentity{&docker.Provider{}}); !errors.Is(err, providercontract.ErrContract) { + t.Fatal("provider registration accepted empty observation identity", err) + } +} diff --git a/services/core/internal/store/provider_operations_fixture_test.go b/services/core/internal/store/provider_operations_fixture_test.go index a7cae2a14..68634162e 100644 --- a/services/core/internal/store/provider_operations_fixture_test.go +++ b/services/core/internal/store/provider_operations_fixture_test.go @@ -9,24 +9,26 @@ import ( 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"}, + "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"}, + "ObservationProviderType": {State: providercontract.Supported}, + "ResolveObservationSource": {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"}, } } func (*lifecycleProvider) Initial(context.Context, sandbox.Reference) (sandbox.Compute, error) { @@ -70,23 +72,35 @@ func (*lifecycleProvider) VerifyCredential(context.Context, []sandbox.Reference) } 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"}, + "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}, + "ObservationProviderType": {State: providercontract.Supported}, + "ResolveObservationSource": {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"}, } } + +func (*lifecycleProvider) ObservationProviderType() string { return "fixture" } +func (p *lifecycleProvider) ResolveObservationSource(context.Context) (runtimeobs.Source, error) { + return p, nil +} + +func (*fakeCheckpointProvider) ObservationProviderType() string { return "fixture" } +func (p *fakeCheckpointProvider) ResolveObservationSource(context.Context) (runtimeobs.Source, error) { + return p, nil +}