Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions contracts/agents-api/runtime-observability.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
9 changes: 6 additions & 3 deletions services/core/cmd/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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 {
Expand Down
17 changes: 13 additions & 4 deletions services/core/cmd/server/managed_generations.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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)
}
Expand All @@ -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)
}

Expand All @@ -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"}
Expand Down
37 changes: 4 additions & 33 deletions services/core/cmd/server/managed_setup.go
Original file line number Diff line number Diff line change
Expand Up @@ -134,51 +134,22 @@ 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
}
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) {
Expand Down
45 changes: 43 additions & 2 deletions services/core/cmd/server/managed_setup_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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"})
Expand All @@ -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")
}
}
Expand Down Expand Up @@ -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 := &microsandbox.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)
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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
}
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Loading
Loading