From 201bbe98afba4461e847f45cddd7ce973ece6973 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Thu, 8 Oct 2026 15:00:06 +0000 Subject: [PATCH 1/2] Declare Runtime observation states once --- contracts/agents-api/core.openapi.yaml | 117 +++++++++++------- .../agents-api/v1/runtime_observations.go | 72 ++++++++--- .../agents-client/src/generated/core-api.ts | 18 +-- packages/agents-client/src/sandbox-client.ts | 4 +- .../core/internal/api/admin_runtime_test.go | 7 +- .../core/internal/api/runtime_observations.go | 28 ++--- .../internal/api/runtime_observations_test.go | 21 ++-- .../core/internal/deployment/allocation.go | 38 +++--- .../core/internal/deployment/allocations.go | 6 +- .../internal/deployment/allocations_test.go | 5 +- .../core/internal/deployment/observation.go | 9 +- .../internal/deployment/observation_test.go | 5 +- services/core/internal/deployment/storage.go | 2 +- .../internal/execution/runtime_observation.go | 14 +-- .../postgres/deploymentpg/allocations.go | 8 +- .../postgres/runtimehistorypg/aggregate.go | 6 +- .../runtimehistorypg/aggregate_test.go | 18 +-- .../postgres/runtimehistorypg/exporter.go | 7 +- .../postgres/runtimehistorypg/reader_test.go | 9 +- .../postgres/runtimehistorypg/samples.go | 5 +- .../postgres/runtimehistorypg/samples_test.go | 5 +- services/core/internal/runtimeobs/exporter.go | 6 +- .../runtimeobs/exporter_independence_test.go | 4 +- services/core/internal/runtimeobs/identity.go | 4 +- .../internal/runtimeobs/operations_test.go | 9 +- .../runtimeobs/otlpexporter/exporter.go | 10 +- .../runtimeobs/otlpexporter/exporter_test.go | 10 +- services/core/internal/runtimeobs/sample.go | 16 --- .../core/internal/runtimeobs/sampler_test.go | 8 +- services/core/internal/runtimeobs/service.go | 27 ++-- .../core/internal/runtimeobs/service_test.go | 76 ++++++------ .../core/internal/runtimeobs/source_test.go | 7 +- .../internal/sandbox/docker/provider_test.go | 3 +- .../core/internal/sandbox/docker/resources.go | 3 +- .../internal/sandbox/docker/resources_test.go | 3 +- .../core/internal/sandbox/e2b/observations.go | 3 +- .../internal/sandbox/e2b/observations_test.go | 3 +- .../sandbox/microsandbox/resources.go | 3 +- .../sandbox/microsandbox/resources_test.go | 3 +- .../internal/sandbox/node/observations.go | 3 +- .../sandbox/node/observations_test.go | 8 +- 41 files changed, 342 insertions(+), 271 deletions(-) diff --git a/contracts/agents-api/core.openapi.yaml b/contracts/agents-api/core.openapi.yaml index d7435b8ed..3a00068fd 100644 --- a/contracts/agents-api/core.openapi.yaml +++ b/contracts/agents-api/core.openapi.yaml @@ -87,24 +87,15 @@ definitions: instance: $ref: '#/definitions/v1.RuntimeInstance' lifecycle_state: - enum: - - active - - sleeping - - transitioning - - pending - - stopped - type: string + allOf: + - $ref: '#/definitions/v1.RuntimeObservationLifecycleState' x-nullable: true memory: allOf: - $ref: '#/definitions/v1.RuntimeMemoryObservation' x-nullable: true mode: - enum: - - none - - self_hosted - - openai_hosted - type: string + $ref: '#/definitions/v1.RuntimeObservationMode' object: enum: - agent.runtime_observation @@ -131,11 +122,7 @@ definitions: type: integer x-nullable: true status: - enum: - - observed - - unsupported - - unavailable - type: string + $ref: '#/definitions/v1.RuntimeObservationStatus' required: - allocation_created_at - cpu @@ -1113,6 +1100,22 @@ definitions: - nodes_on_other_address - self_hosted_executors type: object + deployment.AllocationDiagnostic: + enum: + - "" + - node_unavailable + - resource_missing + - compute_unconfirmed + - ownership_mismatch + - provider_unavailable + type: string + x-enum-varnames: + - AllocationObserved + - AllocationNodeUnavailable + - AllocationResourceMissing + - AllocationComputeUnconfirmed + - AllocationOwnershipMismatch + - AllocationProviderUnavailable deployment.HostHistory: properties: points: @@ -1237,14 +1240,7 @@ definitions: deployment_generation: type: integer diagnostic: - enum: - - "" - - node_unavailable - - resource_missing - - compute_unconfirmed - - ownership_mismatch - - provider_unavailable - type: string + $ref: '#/definitions/deployment.AllocationDiagnostic' environment_id: type: string id: @@ -2879,16 +2875,22 @@ definitions: type: string x-nullable: true kind: - enum: - - managed_allocation - - self_hosted_connection - - none - type: string + $ref: '#/definitions/v1.RuntimeInstanceKind' required: - allocation_id - connection_generation - kind type: object + v1.RuntimeInstanceKind: + enum: + - managed_allocation + - self_hosted_connection + - none + type: string + x-enum-varnames: + - RuntimeInstanceManagedAllocation + - RuntimeInstanceSelfHostedConnection + - RuntimeInstanceNone v1.RuntimeMemoryObservation: properties: limit_bytes: @@ -2923,24 +2925,15 @@ definitions: instance: $ref: '#/definitions/v1.RuntimeInstance' lifecycle_state: - enum: - - active - - sleeping - - transitioning - - pending - - stopped - type: string + allOf: + - $ref: '#/definitions/v1.RuntimeObservationLifecycleState' x-nullable: true memory: allOf: - $ref: '#/definitions/v1.RuntimeMemoryObservation' x-nullable: true mode: - enum: - - none - - self_hosted - - openai_hosted - type: string + $ref: '#/definitions/v1.RuntimeObservationMode' object: enum: - agent.runtime_observation @@ -2967,11 +2960,7 @@ definitions: type: integer x-nullable: true status: - enum: - - observed - - unsupported - - unavailable - type: string + $ref: '#/definitions/v1.RuntimeObservationStatus' required: - allocation_created_at - cpu @@ -2990,6 +2979,40 @@ definitions: - started_at - status type: object + v1.RuntimeObservationLifecycleState: + enum: + - active + - sleeping + - transitioning + - pending + - stopped + type: string + x-enum-varnames: + - RuntimeLifecycleActive + - RuntimeLifecycleSleeping + - RuntimeLifecycleTransitioning + - RuntimeLifecyclePending + - RuntimeLifecycleStopped + v1.RuntimeObservationMode: + enum: + - none + - self_hosted + - openai_hosted + type: string + x-enum-varnames: + - RuntimeModeNone + - RuntimeModeSelfHosted + - RuntimeModeManaged + v1.RuntimeObservationStatus: + enum: + - observed + - unsupported + - unavailable + type: string + x-enum-varnames: + - RuntimeStatusObserved + - RuntimeStatusUnsupported + - RuntimeStatusUnavailable v1.SavedAgent: properties: created_at: diff --git a/contracts/agents-api/v1/runtime_observations.go b/contracts/agents-api/v1/runtime_observations.go index 0a0f46014..6b7b8f570 100644 --- a/contracts/agents-api/v1/runtime_observations.go +++ b/contracts/agents-api/v1/runtime_observations.go @@ -1,28 +1,62 @@ package v1 +type RuntimeObservationMode string + +const ( + RuntimeModeNone RuntimeObservationMode = "none" + RuntimeModeSelfHosted RuntimeObservationMode = "self_hosted" + RuntimeModeManaged RuntimeObservationMode = "openai_hosted" +) + +type RuntimeObservationStatus string + +const ( + RuntimeStatusObserved RuntimeObservationStatus = "observed" + RuntimeStatusUnsupported RuntimeObservationStatus = "unsupported" + RuntimeStatusUnavailable RuntimeObservationStatus = "unavailable" +) + +type RuntimeObservationLifecycleState string + +const ( + RuntimeLifecycleActive RuntimeObservationLifecycleState = "active" + RuntimeLifecycleSleeping RuntimeObservationLifecycleState = "sleeping" + RuntimeLifecycleTransitioning RuntimeObservationLifecycleState = "transitioning" + RuntimeLifecyclePending RuntimeObservationLifecycleState = "pending" + RuntimeLifecycleStopped RuntimeObservationLifecycleState = "stopped" +) + +type RuntimeInstanceKind string + +const ( + RuntimeInstanceManagedAllocation RuntimeInstanceKind = "managed_allocation" + RuntimeInstanceSelfHostedConnection RuntimeInstanceKind = "self_hosted_connection" + RuntimeInstanceNone RuntimeInstanceKind = "none" +) + type RuntimeObservation struct { - ID string `json:"id" binding:"required" format:"uuid"` - Object string `json:"object" enums:"agent.runtime_observation" binding:"required"` - SessionID string `json:"session_id" binding:"required" format:"uuid"` - EnvironmentID *string `json:"environment_id" extensions:"x-nullable" binding:"required" format:"uuid"` - Mode string `json:"mode" enums:"none,self_hosted,openai_hosted" binding:"required"` - ProviderType *string `json:"provider_type" extensions:"x-nullable" binding:"required" pattern:"^[a-z][a-z0-9_]{0,31}$"` - Instance RuntimeInstance `json:"instance" binding:"required"` - LifecycleState *string `json:"lifecycle_state" extensions:"x-nullable" binding:"required" enums:"active,sleeping,transitioning,pending,stopped"` - Status string `json:"status" enums:"observed,unsupported,unavailable" binding:"required"` - Reason *string `json:"reason" extensions:"x-nullable" binding:"required" pattern:"^[a-z][a-z0-9_]{0,95}(?![\\s\\S])"` - AllocationCreatedAt *int64 `json:"allocation_created_at" extensions:"x-nullable" binding:"required" minimum:"0"` - ResolvedAt int64 `json:"resolved_at" binding:"required" minimum:"0"` - ObservedAt *int64 `json:"observed_at" extensions:"x-nullable" binding:"required" minimum:"0"` - StartedAt *int64 `json:"started_at" extensions:"x-nullable" binding:"required" minimum:"0"` - CPU *RuntimeCPUObservation `json:"cpu" extensions:"x-nullable" binding:"required"` - Memory *RuntimeMemoryObservation `json:"memory" extensions:"x-nullable" binding:"required"` + ID string `json:"id" binding:"required" format:"uuid"` + Object string `json:"object" enums:"agent.runtime_observation" binding:"required"` + SessionID string `json:"session_id" binding:"required" format:"uuid"` + EnvironmentID *string `json:"environment_id" extensions:"x-nullable" binding:"required" format:"uuid"` + Mode RuntimeObservationMode `json:"mode" binding:"required"` + ProviderType *string `json:"provider_type" extensions:"x-nullable" binding:"required" pattern:"^[a-z][a-z0-9_]{0,31}$"` + Instance RuntimeInstance `json:"instance" binding:"required"` + LifecycleState *RuntimeObservationLifecycleState `json:"lifecycle_state" extensions:"x-nullable" binding:"required"` + Status RuntimeObservationStatus `json:"status" binding:"required"` + Reason *string `json:"reason" extensions:"x-nullable" binding:"required" pattern:"^[a-z][a-z0-9_]{0,95}(?![\\s\\S])"` + AllocationCreatedAt *int64 `json:"allocation_created_at" extensions:"x-nullable" binding:"required" minimum:"0"` + ResolvedAt int64 `json:"resolved_at" binding:"required" minimum:"0"` + ObservedAt *int64 `json:"observed_at" extensions:"x-nullable" binding:"required" minimum:"0"` + StartedAt *int64 `json:"started_at" extensions:"x-nullable" binding:"required" minimum:"0"` + CPU *RuntimeCPUObservation `json:"cpu" extensions:"x-nullable" binding:"required"` + Memory *RuntimeMemoryObservation `json:"memory" extensions:"x-nullable" binding:"required"` } type RuntimeInstance struct { - Kind string `json:"kind" enums:"managed_allocation,self_hosted_connection,none" binding:"required"` - AllocationID *string `json:"allocation_id" extensions:"x-nullable" binding:"required" format:"uuid"` - ConnectionGeneration *string `json:"connection_generation" extensions:"x-nullable" binding:"required" format:"uuid"` + Kind RuntimeInstanceKind `json:"kind" binding:"required"` + AllocationID *string `json:"allocation_id" extensions:"x-nullable" binding:"required" format:"uuid"` + ConnectionGeneration *string `json:"connection_generation" extensions:"x-nullable" binding:"required" format:"uuid"` } type RuntimeCPUObservation struct { diff --git a/packages/agents-client/src/generated/core-api.ts b/packages/agents-client/src/generated/core-api.ts index 376141492..5ce1a7918 100644 --- a/packages/agents-client/src/generated/core-api.ts +++ b/packages/agents-client/src/generated/core-api.ts @@ -31,9 +31,9 @@ export interface AdminRuntimeObservationDetail { environment_id: string | null; id: string; instance: RuntimeInstance; - lifecycle_state: AdminRuntimeObservationDetailLifecycleState | null; + lifecycle_state: RuntimeObservationLifecycleState | null; memory: RuntimeMemoryObservation | null; - mode: AdminRuntimeObservationDetailMode; + mode: RuntimeObservationMode; object: "agent.runtime_observation"; observed_at: number | null; provider_type: string | null; @@ -41,15 +41,9 @@ export interface AdminRuntimeObservationDetail { resolved_at: number; session_id: string; started_at: number | null; - status: AdminRuntimeObservationDetailStatus; + status: RuntimeObservationStatus; } export const adminRuntimeObservationDetailFields = ["allocation_created_at", "cpu", "disk", "environment_id", "id", "instance", "lifecycle_state", "memory", "mode", "object", "observed_at", "provider_type", "reason", "resolved_at", "session_id", "started_at", "status"] as const; -export const adminRuntimeObservationDetailLifecycleStateValues = ["active", "sleeping", "transitioning", "pending", "stopped"] as const; -export type AdminRuntimeObservationDetailLifecycleState = (typeof adminRuntimeObservationDetailLifecycleStateValues)[number]; -export const adminRuntimeObservationDetailModeValues = ["none", "self_hosted", "openai_hosted"] as const; -export type AdminRuntimeObservationDetailMode = (typeof adminRuntimeObservationDetailModeValues)[number]; -export const adminRuntimeObservationDetailStatusValues = ["observed", "unsupported", "unavailable"] as const; -export type AdminRuntimeObservationDetailStatus = (typeof adminRuntimeObservationDetailStatusValues)[number]; export interface AdminRuntimeObservationList { data: AdminRuntimeObservation[]; first_id: string | null; @@ -112,6 +106,8 @@ export interface AdminauditPage { next_cursor: string; } export const adminauditPageFields = ["data", "has_more", "next_cursor"] as const; +export const allocationDiagnosticValues = ["", "node_unavailable", "resource_missing", "compute_unconfirmed", "ownership_mismatch", "provider_unavailable"] as const; +export type AllocationDiagnostic = (typeof allocationDiagnosticValues)[number]; export interface ConfigurationDiscoveryInput { configuration?: Record; credential?: Record; @@ -429,7 +425,7 @@ export interface NodeAllocation { compute_phase_changed_at: string | null; created_at: string; deployment_generation: number; - diagnostic: NodeAllocationDiagnostic; + diagnostic: AllocationDiagnostic; environment_id: string; id: string; initialization: string; @@ -439,8 +435,6 @@ export interface NodeAllocation { tenant_id: string; } export const nodeAllocationFields = ["compute_phase", "compute_phase_changed_at", "created_at", "deployment_generation", "diagnostic", "environment_id", "id", "initialization", "node_id", "session_id", "state", "tenant_id"] as const; -export const nodeAllocationDiagnosticValues = ["", "node_unavailable", "resource_missing", "compute_unconfirmed", "ownership_mismatch", "provider_unavailable"] as const; -export type NodeAllocationDiagnostic = (typeof nodeAllocationDiagnosticValues)[number]; export interface NodeDetail { active: number; available_disk_bytes: number | null; diff --git a/packages/agents-client/src/sandbox-client.ts b/packages/agents-client/src/sandbox-client.ts index 76eb983b9..1ace80953 100644 --- a/packages/agents-client/src/sandbox-client.ts +++ b/packages/agents-client/src/sandbox-client.ts @@ -4,7 +4,7 @@ import { deploymentContract } from "./deployment-contract"; import { hasOwn, isNonnegativeInteger, isOneOf, isRecord, onlyFields, sameResourceId, schemaFields } from "./response-projection"; import { deploymentResourcesFields, deploymentSpecFields, deploymentSpecRequired, deploymentViewFields, deploymentModeValues, deploymentViewRequired, - hostHistoryFields, hostHistoryPointFields, nodeAllocationDiagnosticValues, nodeAllocationFields, nodeDetailFields, nodeDetailRequired, + hostHistoryFields, hostHistoryPointFields, allocationDiagnosticValues, nodeAllocationFields, nodeDetailFields, nodeDetailRequired, nodeDiagnosticCodeValues, nodeFields, nodeHostFields, nodeRequired, nodeRolloutFields, nodeRolloutRequired, nodeRolloutStateValues, resetModeValues, resetFields, resetOfflineNodeFields, resetRemainingFields, rolloutFields, rolloutNodesFields, rolloutStateValues, runtimeReleaseFields, sandboxAllocationListFields, sandboxNodeListFields, sandboxResourcesFields, sandboxResourcesRequired, suspensionFields, @@ -218,7 +218,7 @@ function projectAllocations(value: unknown, nodeId: string): { data: SandboxAllo valid(Array.isArray(list.data)); return { data: (list.data as unknown[]).map((entry) => { const allocation = members(entry, nodeAllocationFields); - valid(isNonnegativeInteger(allocation.deployment_generation) && strings(allocation, allocationStrings) && sameResourceId(allocation.node_id as string, nodeId) && isOneOf(nodeAllocationDiagnosticValues, allocation.diagnostic) && + valid(isNonnegativeInteger(allocation.deployment_generation) && strings(allocation, allocationStrings) && sameResourceId(allocation.node_id as string, nodeId) && isOneOf(allocationDiagnosticValues, allocation.diagnostic) && nullable(timestamp)(allocation.compute_phase_changed_at) && timestamp(allocation.created_at)); return { ...allocation } as unknown as SandboxAllocation; }) }; diff --git a/services/core/internal/api/admin_runtime_test.go b/services/core/internal/api/admin_runtime_test.go index d2962247a..66c84c42c 100644 --- a/services/core/internal/api/admin_runtime_test.go +++ b/services/core/internal/api/admin_runtime_test.go @@ -9,6 +9,7 @@ import ( "testing" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/projects" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" @@ -65,8 +66,8 @@ func TestAdminRuntimeRoutesRejectHead(t *testing.T) { func unsupportedObservation(session string, at time.Time) runtimeobs.Observation { return runtimeobs.Observation{ - Target: runtimeobs.Target{SessionID: session, Mode: runtimeobs.ModeNone}, - Status: runtimeobs.StatusUnsupported, Reason: "runtime_mode_not_observable", ResolvedAt: at, + Target: runtimeobs.Target{SessionID: session, Mode: v1.RuntimeModeNone}, + Status: v1.RuntimeStatusUnsupported, Reason: "runtime_mode_not_observable", ResolvedAt: at, } } @@ -130,7 +131,7 @@ func TestAdminRuntimeObservationListRejectsWholePageOnIntegrityFailure(t *testin valid, invalid := uuid.NewString(), uuid.NewString() service := runtimeObservationServiceFunc(func(_ context.Context, _, session string) (runtimeobs.Observation, error) { if session == invalid { - return runtimeobs.Observation{Target: runtimeobs.Target{Mode: runtimeobs.ModeNone}, Status: runtimeobs.StatusUnsupported, ResolvedAt: now}, nil + return runtimeobs.Observation{Target: runtimeobs.Target{Mode: v1.RuntimeModeNone}, Status: v1.RuntimeStatusUnsupported, ResolvedAt: now}, nil } return unsupportedObservation(session, now), nil }) diff --git a/services/core/internal/api/runtime_observations.go b/services/core/internal/api/runtime_observations.go index 90ce99da9..e644408ea 100644 --- a/services/core/internal/api/runtime_observations.go +++ b/services/core/internal/api/runtime_observations.go @@ -78,7 +78,7 @@ func runtimeObservationResponse(observation runtimeobs.Observation) (v1.RuntimeO } result := v1.RuntimeObservation{ ID: observation.Target.SessionID, Object: "agent.runtime_observation", SessionID: observation.Target.SessionID, - Mode: string(observation.Target.Mode), Status: string(observation.Status), ResolvedAt: observation.ResolvedAt.Unix(), + Mode: observation.Target.Mode, Status: observation.Status, ResolvedAt: observation.ResolvedAt.Unix(), } if observation.Target.EnvironmentID != "" { result.EnvironmentID = &observation.Target.EnvironmentID @@ -93,8 +93,8 @@ func runtimeObservationResponse(observation runtimeobs.Observation) (v1.RuntimeO result.Reason = &observation.Reason } switch observation.Target.Mode { - case runtimeobs.ModeManaged: - result.Instance.Kind = "managed_allocation" + case v1.RuntimeModeManaged: + result.Instance.Kind = v1.RuntimeInstanceManagedAllocation lifecycleState, err := runtimeLifecycleState(observation.Target.Instance) if err != nil { return v1.RuntimeObservation{}, err @@ -107,13 +107,13 @@ func runtimeObservationResponse(observation runtimeobs.Observation) (v1.RuntimeO created := observation.Target.Instance.AllocationCreatedAt.Unix() result.AllocationCreatedAt = &created } - case runtimeobs.ModeSelfHosted: - result.Instance.Kind = "self_hosted_connection" + case v1.RuntimeModeSelfHosted: + result.Instance.Kind = v1.RuntimeInstanceSelfHostedConnection if observation.Target.Instance.ConnectionGeneration != "" { result.Instance.ConnectionGeneration = &observation.Target.Instance.ConnectionGeneration } - case runtimeobs.ModeNone: - result.Instance.Kind = "none" + case v1.RuntimeModeNone: + result.Instance.Kind = v1.RuntimeInstanceNone default: return v1.RuntimeObservation{}, errors.New("invalid Runtime observation mode") } @@ -135,25 +135,25 @@ func runtimeObservationResponse(observation runtimeobs.Observation) (v1.RuntimeO return result, nil } -func runtimeLifecycleState(instance runtimeobs.Instance) (string, error) { +func runtimeLifecycleState(instance runtimeobs.Instance) (v1.RuntimeObservationLifecycleState, error) { switch instance.AllocationState { case "": if instance.AllocationID == "" { - return "pending", nil + return v1.RuntimeLifecyclePending, nil } return "", errors.New("invalid Runtime allocation state") case "creating": - return "pending", nil + return v1.RuntimeLifecyclePending, nil case "cleanup_pending", "released": - return "stopped", nil + return v1.RuntimeLifecycleStopped, nil case "running": switch instance.ComputePhase { case "suspended": - return "sleeping", nil + return v1.RuntimeLifecycleSleeping, nil case "quiescing", "suspending", "restoring", "waking": - return "transitioning", nil + return v1.RuntimeLifecycleTransitioning, nil case "disabled", "running": - return "active", nil + return v1.RuntimeLifecycleActive, nil default: return "", errors.New("invalid Runtime compute phase") } diff --git a/services/core/internal/api/runtime_observations_test.go b/services/core/internal/api/runtime_observations_test.go index 4a9546419..5e97cd70f 100644 --- a/services/core/internal/api/runtime_observations_test.go +++ b/services/core/internal/api/runtime_observations_test.go @@ -79,7 +79,7 @@ func TestRuntimeObservationRoutesUseSessionIdentityAndExactNullability(t *testin now := time.Date(2026, 9, 22, 8, 0, 0, 0, time.UTC) sessionID := uuid.NewString() service := runtimeObservationFixture{values: map[string]runtimeobs.Observation{ - sessionID: {Target: runtimeobs.Target{SessionID: sessionID, Mode: runtimeobs.ModeNone}, Status: runtimeobs.StatusUnsupported, Reason: "runtime_mode_not_observable", ResolvedAt: now}, + sessionID: {Target: runtimeobs.Target{SessionID: sessionID, Mode: v1.RuntimeModeNone}, Status: v1.RuntimeStatusUnsupported, Reason: "runtime_mode_not_observable", ResolvedAt: now}, }} handler, _, _ := adminTestHandler(t, observeWith(service)) @@ -97,7 +97,7 @@ func TestRuntimeObservationRoutesRejectQueries(t *testing.T) { sessionID := uuid.NewString() now := time.Now().UTC() service := runtimeObservationFixture{values: map[string]runtimeobs.Observation{ - sessionID: {Target: runtimeobs.Target{SessionID: sessionID, Mode: runtimeobs.ModeNone}, Status: runtimeobs.StatusUnsupported, Reason: "runtime_mode_not_observable", ResolvedAt: now}, + sessionID: {Target: runtimeobs.Target{SessionID: sessionID, Mode: v1.RuntimeModeNone}, Status: v1.RuntimeStatusUnsupported, Reason: "runtime_mode_not_observable", ResolvedAt: now}, }} handler, _, _ := adminTestHandler(t, observeWith(service)) invalid := runtimeObservationRequest(handler, adminSessionsPath+sessionID+"/runtime-observation?provider=docker") @@ -112,8 +112,8 @@ func TestRuntimeObservationResponsePreservesObservedZero(t *testing.T) { now := time.Date(2026, 9, 22, 8, 0, 0, 0, time.UTC) sessionID, environmentID := uuid.NewString(), uuid.NewString() value, err := runtimeObservationResponse(runtimeobs.Observation{ - Target: runtimeobs.Target{SessionID: sessionID, EnvironmentID: environmentID, Mode: runtimeobs.ModeManaged, Instance: runtimeobs.Instance{AllocationID: uuid.NewString(), AllocationState: "running", ComputePhase: "running", AllocationCreatedAt: now.Add(-time.Hour)}}, - Status: runtimeobs.StatusObserved, ProviderType: "docker", ResolvedAt: now, + Target: runtimeobs.Target{SessionID: sessionID, EnvironmentID: environmentID, Mode: v1.RuntimeModeManaged, Instance: runtimeobs.Instance{AllocationID: uuid.NewString(), AllocationState: "running", ComputePhase: "running", AllocationCreatedAt: now.Add(-time.Hour)}}, + Status: v1.RuntimeStatusObserved, ProviderType: "docker", ResolvedAt: now, Sample: &runtimeobs.Sample{ObservedAt: now, CPUUsageSecondsTotal: &zeroCPU, MemoryUsageBytes: &zeroMemory}, }) if err != nil || value.LifecycleState == nil || *value.LifecycleState != "active" || value.CPU == nil || value.CPU.UsageSecondsTotal == nil || *value.CPU.UsageSecondsTotal != 0 || value.Memory == nil || value.Memory.UsageBytes == nil || *value.Memory.UsageBytes != 0 { @@ -123,7 +123,8 @@ func TestRuntimeObservationResponsePreservesObservedZero(t *testing.T) { func TestRuntimeLifecycleStateProjectsProviderNeutralPhases(t *testing.T) { for _, item := range []struct { - state, phase, want string + state, phase string + want v1.RuntimeObservationLifecycleState }{ {state: "", phase: "", want: "pending"}, {state: "creating", phase: "disabled", want: "pending"}, @@ -149,10 +150,10 @@ func TestRuntimeObservationResponseRejectsTimesOutsidePublicContract(t *testing. preEpoch := time.Unix(-1, 0).UTC() base := runtimeobs.Observation{ Target: runtimeobs.Target{ - SessionID: uuid.NewString(), EnvironmentID: uuid.NewString(), Mode: runtimeobs.ModeManaged, + SessionID: uuid.NewString(), EnvironmentID: uuid.NewString(), Mode: v1.RuntimeModeManaged, Instance: runtimeobs.Instance{AllocationID: uuid.NewString(), AllocationCreatedAt: now.Add(-time.Hour)}, }, - Status: runtimeobs.StatusObserved, ResolvedAt: now, + Status: v1.RuntimeStatusObserved, ResolvedAt: now, Sample: &runtimeobs.Sample{ObservedAt: now, StartedAt: timePointer(now.Add(-time.Minute))}, } for _, mutate := range []func(*runtimeobs.Observation){ @@ -178,9 +179,9 @@ func TestAdminRuntimeObservationProjectsReportedUtilizationAndDisk(t *testing.T) ratio, cores := .1955, 2.0 memoryUsed, memoryTotal, diskUsed, diskTotal := uint64(183836672), uint64(2079141888), uint64(1593188352), uint64(23511863296) observation := runtimeobs.Observation{ - Target: runtimeobs.Target{SessionID: uuid.NewString(), EnvironmentID: uuid.NewString(), Mode: runtimeobs.ModeManaged, + Target: runtimeobs.Target{SessionID: uuid.NewString(), EnvironmentID: uuid.NewString(), Mode: v1.RuntimeModeManaged, Instance: runtimeobs.Instance{AllocationID: uuid.NewString(), AllocationState: "running", ComputePhase: "disabled", AllocationCreatedAt: now.Add(-time.Minute)}}, - Status: runtimeobs.StatusObserved, ProviderType: "e2b", ResolvedAt: now, + Status: v1.RuntimeStatusObserved, ProviderType: "e2b", ResolvedAt: now, Sample: &runtimeobs.Sample{ObservedAt: now, StartedAt: timePointer(now.Add(-time.Minute)), CPUUtilizationRatio: &ratio, CPUCapacityCores: &cores, MemoryUsageBytes: &memoryUsed, MemoryLimitBytes: &memoryTotal, DiskUsageBytes: &diskUsed, DiskLimitBytes: &diskTotal}, } @@ -223,7 +224,7 @@ func (s declaredObservationSource) Observe(context.Context, runtimeobs.Target) ( func (s declaredObservationSource) Resolve(_ context.Context, tenant, session string) (runtimeobs.Target, error) { return runtimeobs.Target{ - TenantID: tenant, SessionID: session, EnvironmentID: "22222222-2222-4222-8222-222222222222", Mode: runtimeobs.ModeManaged, + TenantID: tenant, SessionID: session, EnvironmentID: "22222222-2222-4222-8222-222222222222", Mode: v1.RuntimeModeManaged, Instance: runtimeobs.Instance{AllocationID: "33333333-3333-4333-8333-333333333333", ProviderKey: "provider", AllocationState: "running", ComputePhase: "running"}, }, nil } diff --git a/services/core/internal/deployment/allocation.go b/services/core/internal/deployment/allocation.go index 6cab5cb4e..9aebb8af9 100644 --- a/services/core/internal/deployment/allocation.go +++ b/services/core/internal/deployment/allocation.go @@ -12,7 +12,7 @@ type Allocation struct { DeploymentGeneration uint64 // NodeID is empty for an allocation no node serves. NodeID string - ObservationError string + ObservationError AllocationDiagnostic ComputePhase string ComputeRevision int64 ComputeState json.RawMessage @@ -128,27 +128,33 @@ type LifecyclePlacement struct { // NodeAllocation is an unreleased allocation a node serves. type NodeAllocation struct { - DeploymentGeneration uint64 `json:"deployment_generation" binding:"required"` - Diagnostic string `json:"diagnostic" enums:",node_unavailable,resource_missing,compute_unconfirmed,ownership_mismatch,provider_unavailable" binding:"required"` - ID string `json:"id" binding:"required"` - NodeID string `json:"node_id" binding:"required"` - TenantID string `json:"tenant_id" binding:"required"` - SessionID string `json:"session_id" binding:"required"` - EnvironmentID string `json:"environment_id" binding:"required"` - State string `json:"state" binding:"required"` - ComputePhase string `json:"compute_phase" binding:"required"` + DeploymentGeneration uint64 `json:"deployment_generation" binding:"required"` + Diagnostic AllocationDiagnostic `json:"diagnostic" binding:"required"` + ID string `json:"id" binding:"required"` + NodeID string `json:"node_id" binding:"required"` + TenantID string `json:"tenant_id" binding:"required"` + SessionID string `json:"session_id" binding:"required"` + EnvironmentID string `json:"environment_id" binding:"required"` + State string `json:"state" binding:"required"` + ComputePhase string `json:"compute_phase" binding:"required"` // The time the allocation entered its current compute_phase, or null when unknown; an allocation that existed before Core recorded it reports null until its next phase change. For a suspended microsandbox allocation, this time plus the deployment's snapshot retention tells roughly when Core reclaims it. ComputePhaseChangedAt *time.Time `json:"compute_phase_changed_at" extensions:"x-nullable" binding:"required"` Initialization string `json:"initialization" binding:"required"` CreatedAt time.Time `json:"created_at" binding:"required"` } -// observationDiagnostics are the diagnostics an observation records; empty -// clears the previous one. -var observationDiagnostics = map[string]bool{ - "": true, "node_unavailable": true, "resource_missing": true, - "compute_unconfirmed": true, "ownership_mismatch": true, "provider_unavailable": true, -} +// AllocationDiagnostic is the sanitized result of a node-backed observation. +// The empty value clears the previous diagnostic. +type AllocationDiagnostic string + +const ( + AllocationObserved AllocationDiagnostic = "" + AllocationNodeUnavailable AllocationDiagnostic = "node_unavailable" + AllocationResourceMissing AllocationDiagnostic = "resource_missing" + AllocationComputeUnconfirmed AllocationDiagnostic = "compute_unconfirmed" + AllocationOwnershipMismatch AllocationDiagnostic = "ownership_mismatch" + AllocationProviderUnavailable AllocationDiagnostic = "provider_unavailable" +) // computeTransition reports whether managed compute may move from one phase // to another. Repeating a phase records a new state, except while disabled. diff --git a/services/core/internal/deployment/allocations.go b/services/core/internal/deployment/allocations.go index 48fda2384..2bd0ecbc5 100644 --- a/services/core/internal/deployment/allocations.go +++ b/services/core/internal/deployment/allocations.go @@ -245,8 +245,10 @@ func (e *ExecutionOperations) SetCompute(ctx context.Context, owner Allocation, // RecordObservation records a diagnostic for the observed compute revision // and state of a node-backed allocation. It never releases ownership, changes // public readiness or authorizes replacement. -func (e *ExecutionOperations) RecordObservation(ctx context.Context, owner Allocation, diagnostic string) error { - if !observationDiagnostics[diagnostic] { +func (e *ExecutionOperations) RecordObservation(ctx context.Context, owner Allocation, diagnostic AllocationDiagnostic) error { + switch diagnostic { + case AllocationObserved, AllocationNodeUnavailable, AllocationResourceMissing, AllocationComputeUnconfirmed, AllocationOwnershipMismatch, AllocationProviderUnavailable: + default: return ErrInvalidInput } if owner.NodeID == "" { diff --git a/services/core/internal/deployment/allocations_test.go b/services/core/internal/deployment/allocations_test.go index 00303daf3..b2c98b8ac 100644 --- a/services/core/internal/deployment/allocations_test.go +++ b/services/core/internal/deployment/allocations_test.go @@ -10,10 +10,9 @@ import ( "testing" "time" - "github.com/google/uuid" - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment/placement" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" + "github.com/google/uuid" ) type fakeReservationTx struct { @@ -108,7 +107,7 @@ func (f *fakeAllocationTx) SetCompute(Allocation, ComputeChange) (Allocation, er return Allocation{}, nil } -func (f *fakeAllocationTx) RecordObservation(Allocation, string) error { +func (f *fakeAllocationTx) RecordObservation(Allocation, AllocationDiagnostic) error { unexpected(f.t, "RecordObservation") return nil } diff --git a/services/core/internal/deployment/observation.go b/services/core/internal/deployment/observation.go index cf26da4a2..7fad73a7b 100644 --- a/services/core/internal/deployment/observation.go +++ b/services/core/internal/deployment/observation.go @@ -8,6 +8,7 @@ import ( "fmt" "io" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" ) @@ -39,7 +40,7 @@ func (r *ObservationResolver) Resolve(ctx context.Context, tenantID, sessionID s if err := json.Unmarshal(session.Configuration, &configuration); err != nil || configuration.Environment == nil { return runtimeobs.Target{}, errors.New("invalid stored Runtime environment configuration") } - target := runtimeobs.Target{TenantID: session.TenantID, SessionID: session.ID, Mode: runtimeobs.Mode(configuration.Environment.Type)} + target := runtimeobs.Target{TenantID: session.TenantID, SessionID: session.ID, Mode: v1.RuntimeObservationMode(configuration.Environment.Type)} // Telemetry counts measured usage continuously, including active Turns. // Public Session usage stays null until every root Turn ends measured. measured, err := r.sessions.MeasuredSessionUsage(ctx, session.TenantID, session.ID) @@ -52,12 +53,12 @@ func (r *ObservationResolver) Resolve(ctx context.Context, tenantID, sessionID s } target.TokenUsage = usage switch target.Mode { - case runtimeobs.ModeNone: + case v1.RuntimeModeNone: if session.Environment != nil { return runtimeobs.Target{}, errors.New("environment:none unexpectedly has a durable Environment") } return target, nil - case runtimeobs.ModeSelfHosted: + case v1.RuntimeModeSelfHosted: if session.Environment == nil { return runtimeobs.Target{}, errors.New("self-hosted Session is missing its Environment") } @@ -66,7 +67,7 @@ func (r *ObservationResolver) Resolve(ctx context.Context, tenantID, sessionID s } target.EnvironmentID = session.Environment.ID return target, nil - case runtimeobs.ModeManaged: + case v1.RuntimeModeManaged: if session.Environment == nil { return runtimeobs.Target{}, errors.New("managed Session is missing its Environment") } diff --git a/services/core/internal/deployment/observation_test.go b/services/core/internal/deployment/observation_test.go index 27be18e2c..0954620b0 100644 --- a/services/core/internal/deployment/observation_test.go +++ b/services/core/internal/deployment/observation_test.go @@ -6,6 +6,7 @@ import ( "errors" "testing" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" ) @@ -66,7 +67,7 @@ func TestObservationResolverBindsManagedSessionEnvironmentAndAllocation(t *testi if err != nil { t.Fatal(err) } - if target.TenantID != "tenant" || target.SessionID != "session" || target.EnvironmentID != "environment" || target.Mode != runtimeobs.ModeManaged || target.Instance.AllocationID != "allocation" || target.Instance.ProviderKey != "provider" { + if target.TenantID != "tenant" || target.SessionID != "session" || target.EnvironmentID != "environment" || target.Mode != v1.RuntimeModeManaged || target.Instance.AllocationID != "allocation" || target.Instance.ProviderKey != "provider" { t.Fatalf("incorrect managed identity binding: %+v", target) } if string(target.Instance.ProviderState) != `{"current":{"name":"sandbox"}}` { @@ -165,7 +166,7 @@ func TestObservationResolverReportsManagedAllocationAsUnavailable(t *testing.T) t.Fatal(err) } target, err := r.Resolve(t.Context(), "tenant", "session") - if !errors.Is(err, runtimeobs.ErrUnavailable) || target.EnvironmentID != "environment" || target.Mode != runtimeobs.ModeManaged { + if !errors.Is(err, runtimeobs.ErrUnavailable) || target.EnvironmentID != "environment" || target.Mode != v1.RuntimeModeManaged { t.Fatalf("allocation absence was not preserved: %+v %v", target, err) } } diff --git a/services/core/internal/deployment/storage.go b/services/core/internal/deployment/storage.go index c2ce9d00d..87d42465d 100644 --- a/services/core/internal/deployment/storage.go +++ b/services/core/internal/deployment/storage.go @@ -119,7 +119,7 @@ type AllocationTx interface { SetCompute(current Allocation, change ComputeChange) (Allocation, error) // RecordObservation records the diagnostic for current's compute // revision and state. It changes nothing once either moved on. - RecordObservation(current Allocation, diagnostic string) error + RecordObservation(current Allocation, diagnostic AllocationDiagnostic) error } // AllocationCleanupTx is one Session-locked allocation cleanup. diff --git a/services/core/internal/execution/runtime_observation.go b/services/core/internal/execution/runtime_observation.go index 7d8b0953d..6ae8a72bc 100644 --- a/services/core/internal/execution/runtime_observation.go +++ b/services/core/internal/execution/runtime_observation.go @@ -12,22 +12,22 @@ func (r *runtimeLifecycle) recordObservation(ctx context.Context, owner deployme if owner.NodeID == "" { return } - diagnostic := "" + diagnostic := deployment.AllocationObserved if observed != nil { switch { case errors.Is(observed, sandbox.ErrOwnership): - diagnostic = "ownership_mismatch" + diagnostic = deployment.AllocationOwnershipMismatch case errors.Is(observed, sandbox.ErrNotFound): - diagnostic = "compute_unconfirmed" + diagnostic = deployment.AllocationComputeUnconfirmed if owner.CreateSettled { - diagnostic = "resource_missing" + diagnostic = deployment.AllocationResourceMissing } case errors.Is(observed, sandbox.ErrComputeUnconfirmed): - diagnostic = "compute_unconfirmed" + diagnostic = deployment.AllocationComputeUnconfirmed default: - diagnostic = "provider_unavailable" + diagnostic = deployment.AllocationProviderUnavailable if online, err := r.reader.NodeOnline(ctx, owner.NodeID); err == nil && !online { - diagnostic = "node_unavailable" + diagnostic = deployment.AllocationNodeUnavailable } } } diff --git a/services/core/internal/persistence/postgres/deploymentpg/allocations.go b/services/core/internal/persistence/postgres/deploymentpg/allocations.go index e8b14f73c..38c7e0090 100644 --- a/services/core/internal/persistence/postgres/deploymentpg/allocations.go +++ b/services/core/internal/persistence/postgres/deploymentpg/allocations.go @@ -20,7 +20,7 @@ import ( // allocation converts a stored allocation with the facts of its Session. func allocation(row sqlc.RuntimeAllocation, session, tenant pgtype.UUID, deleted pgtype.Timestamptz, expired bool) deployment.Allocation { return deployment.Allocation{ - DeploymentGeneration: uint64(row.DeploymentGeneration.Int64), NodeID: uuidString(row.NodeID), ObservationError: row.ObservationError, + DeploymentGeneration: uint64(row.DeploymentGeneration.Int64), NodeID: uuidString(row.NodeID), ObservationError: deployment.AllocationDiagnostic(row.ObservationError), ComputePhase: row.ComputePhase, ComputeRevision: row.ComputeRevision, ComputeState: row.ComputeState, ComputeActivityAt: row.ComputeActivityAt.Time, ComputeWakeRequested: row.ComputeWakeRequested, ComputeRetainedUntil: timestamp(row.ComputeRetainedUntil), ID: uuidString(row.ID), EnvironmentID: uuidString(row.EnvironmentID), SessionID: uuidString(session), TenantID: uuidString(tenant), @@ -282,12 +282,12 @@ func (t *allocationTx) SetCompute(current deployment.Allocation, change deployme }) } -func (t *allocationTx) RecordObservation(current deployment.Allocation, diagnostic string) error { +func (t *allocationTx) RecordObservation(current deployment.Allocation, diagnostic deployment.AllocationDiagnostic) error { id, err := parseID(current.ID) if err != nil { return err } - return t.q.SetRuntimeObservation(t.ctx, sqlc.SetRuntimeObservationParams{ID: id, ComputeRevision: current.ComputeRevision, State: current.State, ObservationError: diagnostic}) + return t.q.SetRuntimeObservation(t.ctx, sqlc.SetRuntimeObservationParams{ID: id, ComputeRevision: current.ComputeRevision, State: current.State, ObservationError: string(diagnostic)}) } // change applies a guarded write to current and returns the stored result @@ -401,7 +401,7 @@ func (s *Store) NodeAllocations(ctx context.Context, nodeID string) ([]deploymen result := make([]deployment.NodeAllocation, 0, len(rows)) for _, a := range rows { result = append(result, deployment.NodeAllocation{ - DeploymentGeneration: uint64(a.DeploymentGeneration.Int64), Diagnostic: a.ObservationError, ID: uuidString(a.ID), NodeID: uuidString(a.NodeID), + DeploymentGeneration: uint64(a.DeploymentGeneration.Int64), Diagnostic: deployment.AllocationDiagnostic(a.ObservationError), ID: uuidString(a.ID), NodeID: uuidString(a.NodeID), TenantID: uuidString(a.TenantID), SessionID: uuidString(a.SessionID), EnvironmentID: uuidString(a.EnvironmentID), State: a.State, ComputePhase: a.ComputePhase, ComputePhaseChangedAt: timestamp(a.ComputePhaseChangedAt), Initialization: a.Initialization, CreatedAt: a.CreatedAt.Time, }) diff --git a/services/core/internal/persistence/postgres/runtimehistorypg/aggregate.go b/services/core/internal/persistence/postgres/runtimehistorypg/aggregate.go index f869706e8..74ed73645 100644 --- a/services/core/internal/persistence/postgres/runtimehistorypg/aggregate.go +++ b/services/core/internal/persistence/postgres/runtimehistorypg/aggregate.go @@ -7,8 +7,8 @@ import ( "sort" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimehistory" - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" "github.com/google/uuid" ) @@ -26,7 +26,7 @@ type rawSample struct { allocation string startedAt *time.Time provider string - status runtimeobs.Status + status v1.RuntimeObservationStatus hasSample bool metrics map[string]float64 } @@ -307,7 +307,7 @@ func safeUint64(value float64) (uint64, error) { func addCoverage(value *coverageAggregate, sample *rawSample) { value.observations++ - if sample.status == runtimeobs.StatusObserved { + if sample.status == v1.RuntimeStatusObserved { value.observed++ } else { value.unavailable++ diff --git a/services/core/internal/persistence/postgres/runtimehistorypg/aggregate_test.go b/services/core/internal/persistence/postgres/runtimehistorypg/aggregate_test.go index de23739a7..5eba9017e 100644 --- a/services/core/internal/persistence/postgres/runtimehistorypg/aggregate_test.go +++ b/services/core/internal/persistence/postgres/runtimehistorypg/aggregate_test.go @@ -5,7 +5,7 @@ import ( "testing" "time" - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/google/uuid" ) @@ -18,12 +18,12 @@ func TestAggregateSeriesUsesProviderObservationOrder(t *testing.T) { samples := []*rawSample{ { resolvedAt: secondResolved, observedAt: &secondObserved, allocation: testAllocation, startedAt: &startedAt, - provider: "docker", status: runtimeobs.StatusObserved, hasSample: true, + provider: "docker", status: v1.RuntimeStatusObserved, hasSample: true, metrics: map[string]float64{CPUUsageName: 3, CPUCapacityName: 2, MemoryUsageName: 200}, }, { resolvedAt: firstResolved, observedAt: &firstObserved, allocation: testAllocation, startedAt: &startedAt, - provider: "docker", status: runtimeobs.StatusObserved, hasSample: true, + provider: "docker", status: v1.RuntimeStatusObserved, hasSample: true, metrics: map[string]float64{CPUUsageName: 1, CPUCapacityName: 2, MemoryUsageName: 100}, }, } @@ -44,17 +44,17 @@ func TestAggregateSeriesLeavesCPUUsageGapWhenEitherEndpointLacksCapacity(t *test samples := []*rawSample{ { resolvedAt: firstObserved, observedAt: &firstObserved, allocation: testAllocation, startedAt: &startedAt, - provider: "docker", status: runtimeobs.StatusObserved, hasSample: true, + provider: "docker", status: v1.RuntimeStatusObserved, hasSample: true, metrics: map[string]float64{CPUUsageName: 1, CPUCapacityName: 2}, }, { resolvedAt: secondObserved, observedAt: &secondObserved, allocation: testAllocation, startedAt: &startedAt, - provider: "docker", status: runtimeobs.StatusObserved, hasSample: true, + provider: "docker", status: v1.RuntimeStatusObserved, hasSample: true, metrics: map[string]float64{CPUUsageName: 3}, }, { resolvedAt: thirdObserved, observedAt: &thirdObserved, allocation: testAllocation, startedAt: &startedAt, - provider: "docker", status: runtimeobs.StatusObserved, hasSample: true, + provider: "docker", status: v1.RuntimeStatusObserved, hasSample: true, metrics: map[string]float64{CPUUsageName: 5, CPUCapacityName: 2}, }, } @@ -77,14 +77,14 @@ func TestAggregateIgnoresLookbackOnlySeriesForSeriesLimit(t *testing.T) { observedAt := start.Add(-2 * time.Second) raw = append(raw, &rawSample{ resolvedAt: start.Add(-time.Second), observedAt: &observedAt, allocation: uuid.NewString(), startedAt: &startedAt, - provider: "docker", status: runtimeobs.StatusObserved, hasSample: true, metrics: map[string]float64{}, + provider: "docker", status: v1.RuntimeStatusObserved, hasSample: true, metrics: map[string]float64{}, }) } startedAt := start.Add(-time.Minute) observedAt := start.Add(5 * time.Second) raw = append(raw, &rawSample{ resolvedAt: start.Add(6 * time.Second), observedAt: &observedAt, allocation: uuid.NewString(), startedAt: &startedAt, - provider: "docker", status: runtimeobs.StatusObserved, hasSample: true, metrics: map[string]float64{}, + provider: "docker", status: v1.RuntimeStatusObserved, hasSample: true, metrics: map[string]float64{}, }) result, err := aggregate(query, start.Add(time.Minute), raw) if err != nil { @@ -104,7 +104,7 @@ func TestAggregateSeriesAveragesProviderReportedUtilization(t *testing.T) { observedAt := start.Add(time.Duration(5+index*15) * time.Second) samples = append(samples, &rawSample{ resolvedAt: observedAt, observedAt: &observedAt, allocation: testAllocation, startedAt: &startedAt, - provider: "e2b", status: runtimeobs.StatusObserved, hasSample: true, + provider: "e2b", status: v1.RuntimeStatusObserved, hasSample: true, metrics: map[string]float64{CPUUtilizationName: ratio, CPUCapacityName: 2, MemoryUsageName: 100}, }) } diff --git a/services/core/internal/persistence/postgres/runtimehistorypg/exporter.go b/services/core/internal/persistence/postgres/runtimehistorypg/exporter.go index 13935943c..afcb6e7b5 100644 --- a/services/core/internal/persistence/postgres/runtimehistorypg/exporter.go +++ b/services/core/internal/persistence/postgres/runtimehistorypg/exporter.go @@ -7,6 +7,7 @@ import ( "regexp" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" "github.com/google/uuid" ) @@ -51,13 +52,13 @@ func validateRecord(record runtimeobs.ExportRecord) error { if record.AllocationID != "" && !validID(record.AllocationID) { return invalid } - if record.Mode != runtimeobs.ModeManaged || record.CollectionSource != runtimeobs.CollectionSourcePeriodic || !validTime(record.ResolvedAt) || record.ProviderType != "" && !providerPattern.MatchString(record.ProviderType) { + if record.Mode != v1.RuntimeModeManaged || record.CollectionSource != runtimeobs.CollectionSourcePeriodic || !validTime(record.ResolvedAt) || record.ProviderType != "" && !providerPattern.MatchString(record.ProviderType) { return invalid } - if record.Status != runtimeobs.StatusObserved && record.Status != runtimeobs.StatusUnavailable && record.Status != runtimeobs.StatusUnsupported { + if record.Status != v1.RuntimeStatusObserved && record.Status != v1.RuntimeStatusUnavailable && record.Status != v1.RuntimeStatusUnsupported { return invalid } - if (record.Status == runtimeobs.StatusObserved) != (record.Sample != nil) { + if (record.Status == v1.RuntimeStatusObserved) != (record.Sample != nil) { return invalid } if value := record.Sample; value != nil { diff --git a/services/core/internal/persistence/postgres/runtimehistorypg/reader_test.go b/services/core/internal/persistence/postgres/runtimehistorypg/reader_test.go index 85bc4c916..8a1e388eb 100644 --- a/services/core/internal/persistence/postgres/runtimehistorypg/reader_test.go +++ b/services/core/internal/persistence/postgres/runtimehistorypg/reader_test.go @@ -7,6 +7,7 @@ import ( "testing" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimehistory" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" ) @@ -21,7 +22,7 @@ const ( func testCapabilities() runtimehistory.Capabilities { return runtimehistory.Capabilities{ SampleInterval: 30 * time.Second, - Retention: 7 * 24 * time.Hour, MinimumStep: 30 * time.Second, MaximumRange: 24 * time.Hour, + Retention: 7 * 24 * time.Hour, MinimumStep: 30 * time.Second, MaximumRange: 24 * time.Hour, MaximumPoints: 1_000, MaximumSeries: 64, MaximumTotalPoints: 10_000, Metrics: []runtimehistory.Metric{runtimehistory.MetricCPU, runtimehistory.MetricMemory}, } @@ -76,7 +77,7 @@ func observedRecord(at, started time.Time, cpu float64) runtimeobs.ExportRecord capacity := 2.0 memory := uint64(123) return runtimeobs.ExportRecord{TenantID: testTenant, SessionID: testSession, EnvironmentID: testEnvironment, AllocationID: testAllocation, - Mode: runtimeobs.ModeManaged, ProviderType: "docker", Status: runtimeobs.StatusObserved, CollectionSource: runtimeobs.CollectionSourcePeriodic, + Mode: v1.RuntimeModeManaged, ProviderType: "docker", Status: v1.RuntimeStatusObserved, CollectionSource: runtimeobs.CollectionSourcePeriodic, ResolvedAt: at, Sample: &runtimeobs.Sample{ObservedAt: at, StartedAt: &started, CPUUsageSecondsTotal: &cpu, CPUCapacityCores: &capacity, MemoryUsageBytes: &memory}, TokenUsage: &runtimeobs.TokenUsage{InputTokens: 12, OutputTokens: 3}} } @@ -151,7 +152,7 @@ func TestExporterPersistsOnlyPeriodicAndPreservesUnknownMetrics(t *testing.T) { t.Fatal(s.written) } record.Sample = nil - record.Status = runtimeobs.StatusUnavailable + record.Status = v1.RuntimeStatusUnavailable if err := r.Export(t.Context(), record); err != nil { t.Fatal(err) } @@ -199,7 +200,7 @@ func TestUnavailableObservationBreaksCPUContinuity(t *testing.T) { started := start.Add(-time.Hour) unavailable := observedRecord(start.Add(20*time.Second), started, 0) unavailable.Sample = nil - unavailable.Status = runtimeobs.StatusUnavailable + unavailable.Status = v1.RuntimeStatusUnavailable s := &fakeStore{records: []runtimeobs.ExportRecord{observedRecord(start.Add(5*time.Second), started, 1), unavailable, observedRecord(start.Add(35*time.Second), started, 100)}} r := testReader(t, s, start.Add(time.Minute)) result, err := r.Query(t.Context(), testQuery(start)) diff --git a/services/core/internal/persistence/postgres/runtimehistorypg/samples.go b/services/core/internal/persistence/postgres/runtimehistorypg/samples.go index 08f637f8f..f83ebbbf1 100644 --- a/services/core/internal/persistence/postgres/runtimehistorypg/samples.go +++ b/services/core/internal/persistence/postgres/runtimehistorypg/samples.go @@ -4,6 +4,7 @@ import ( "context" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/db/sqlc" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" @@ -83,8 +84,8 @@ func (s samples) ListRuntimeHistorySamples(ctx context.Context, tenantID, sessio for _, row := range rows { record := runtimeobs.ExportRecord{ TenantID: tenantID, SessionID: sessionID, EnvironmentID: environmentID, - Mode: runtimeobs.ModeManaged, CollectionSource: runtimeobs.CollectionSourcePeriodic, - ResolvedAt: time.Unix(0, row.ResolvedAtNs).UTC(), ProviderType: row.ProviderType, Status: runtimeobs.Status(row.Status), + Mode: v1.RuntimeModeManaged, CollectionSource: runtimeobs.CollectionSourcePeriodic, + ResolvedAt: time.Unix(0, row.ResolvedAtNs).UTC(), ProviderType: row.ProviderType, Status: v1.RuntimeObservationStatus(row.Status), } if row.AllocationID.Valid { record.AllocationID = uuid.UUID(row.AllocationID.Bytes).String() diff --git a/services/core/internal/persistence/postgres/runtimehistorypg/samples_test.go b/services/core/internal/persistence/postgres/runtimehistorypg/samples_test.go index aae9979be..d4a5ddabc 100644 --- a/services/core/internal/persistence/postgres/runtimehistorypg/samples_test.go +++ b/services/core/internal/persistence/postgres/runtimehistorypg/samples_test.go @@ -5,6 +5,7 @@ import ( "testing" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimehistory" @@ -42,7 +43,7 @@ func historyOwner(t *testing.T, pool *pgxpool.Pool) runtimehistory.Scope { func historyRecord(scope runtimehistory.Scope, allocation string, started, at time.Time, cpu float64, input uint64) runtimeobs.ExportRecord { capacity := 2.0 memory := uint64(512) - return runtimeobs.ExportRecord{TenantID: scope.TenantID, SessionID: scope.SessionID, EnvironmentID: scope.EnvironmentID, AllocationID: allocation, Mode: runtimeobs.ModeManaged, ProviderType: "docker", Status: runtimeobs.StatusObserved, CollectionSource: runtimeobs.CollectionSourcePeriodic, ResolvedAt: at, + return runtimeobs.ExportRecord{TenantID: scope.TenantID, SessionID: scope.SessionID, EnvironmentID: scope.EnvironmentID, AllocationID: allocation, Mode: v1.RuntimeModeManaged, ProviderType: "docker", Status: v1.RuntimeStatusObserved, CollectionSource: runtimeobs.CollectionSourcePeriodic, ResolvedAt: at, Sample: &runtimeobs.Sample{ObservedAt: at, StartedAt: &started, CPUUsageSecondsTotal: &cpu, CPUCapacityCores: &capacity, MemoryUsageBytes: &memory}, TokenUsage: &runtimeobs.TokenUsage{InputTokens: input, OutputTokens: input / 2}} } func historyQuery(scope runtimehistory.Scope, start, end time.Time, points int) runtimehistory.Query { @@ -68,7 +69,7 @@ func TestPostgresRuntimeHistoryAcceptance(t *testing.T) { historyRecord(scope, allocation, started, start.Add(95*time.Second), 90, 40), } unavailable := historyRecord(scope, allocation, started, start.Add(155*time.Second), 100, 50) - unavailable.Status = runtimeobs.StatusUnavailable + unavailable.Status = v1.RuntimeStatusUnavailable unavailable.Sample = nil records = append(records, unavailable) for _, record := range records { diff --git a/services/core/internal/runtimeobs/exporter.go b/services/core/internal/runtimeobs/exporter.go index 547b6bdd9..6f0533448 100644 --- a/services/core/internal/runtimeobs/exporter.go +++ b/services/core/internal/runtimeobs/exporter.go @@ -5,6 +5,8 @@ import ( "errors" "sync" "time" + + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" ) const ( @@ -22,9 +24,9 @@ type ExportRecord struct { EnvironmentID string AllocationID string - Mode Mode + Mode v1.RuntimeObservationMode ProviderType string - Status Status + Status v1.RuntimeObservationStatus Reason string CollectionSource CollectionSource diff --git a/services/core/internal/runtimeobs/exporter_independence_test.go b/services/core/internal/runtimeobs/exporter_independence_test.go index 3a24a74ec..0ef735cb7 100644 --- a/services/core/internal/runtimeobs/exporter_independence_test.go +++ b/services/core/internal/runtimeobs/exporter_independence_test.go @@ -3,11 +3,13 @@ package runtimeobs import ( "testing" "time" + + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" ) func TestExporterOutageDoesNotDelayOtherDestinations(t *testing.T) { now := time.Now() - target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} + target := Target{EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, 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}, sourceOf(&fixedSource{sample: Sample{ObservedAt: now}}), diff --git a/services/core/internal/runtimeobs/identity.go b/services/core/internal/runtimeobs/identity.go index bd78c5633..b9b468ebe 100644 --- a/services/core/internal/runtimeobs/identity.go +++ b/services/core/internal/runtimeobs/identity.go @@ -3,6 +3,8 @@ package runtimeobs import ( "encoding/json" "time" + + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" ) // Instance is one provider-owned Runtime incarnation. AllocationID is present @@ -24,7 +26,7 @@ type Instance struct { // process or sandbox, so callers must retain the complete binding. type Target struct { TenantID, SessionID, EnvironmentID string - Mode Mode + Mode v1.RuntimeObservationMode Instance Instance // TokenUsage is the latest measured Core Session usage: every recorded // root Turn snapshot, active Turns included. Unlike public Session usage it diff --git a/services/core/internal/runtimeobs/operations_test.go b/services/core/internal/runtimeobs/operations_test.go index 0fc73a2c0..0d4843557 100644 --- a/services/core/internal/runtimeobs/operations_test.go +++ b/services/core/internal/runtimeobs/operations_test.go @@ -7,6 +7,7 @@ import ( "reflect" "testing" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/providercontract" ) @@ -17,13 +18,13 @@ func (*unsupportedObservation) ProviderOperations() providercontract.Operations } func TestUnsupportedObservationIsNotUnavailable(t *testing.T) { source := &unsupportedObservation{&fixedSource{}} - target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} + target := Target{EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} service, err := NewService(fixedResolver{target: target}, sourceOf(source)) if err != nil { t.Fatal(err) } observation, err := service.ObserveSession(t.Context(), "tenant", "session") - if err != nil || observation.Status != StatusUnsupported || observation.Reason != "native_metrics_not_supported" || source.calls != 0 { + if err != nil || observation.Status != v1.RuntimeStatusUnsupported || observation.Reason != "native_metrics_not_supported" || source.calls != 0 { t.Fatal(observation, err, source.calls) } } @@ -60,7 +61,7 @@ func TestUnavailableReasonsMatchSharedFixture(t *testing.T) { {state: "running", observeErr: context.DeadlineExceeded}, {state: "running", observeErr: ErrUnavailable}, } { - target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: scenario.state}} + target := Target{EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: scenario.state}} if scenario.resolveErr != nil { target.Instance = Instance{} } @@ -69,7 +70,7 @@ func TestUnavailableReasonsMatchSharedFixture(t *testing.T) { t.Fatal(err) } observation, err := service.ObserveSession(t.Context(), "tenant", "session") - if err != nil || observation.Status != StatusUnavailable { + if err != nil || observation.Status != v1.RuntimeStatusUnavailable { t.Fatalf("classification = %+v, %v", observation, err) } actual[observation.Reason] = true diff --git a/services/core/internal/runtimeobs/otlpexporter/exporter.go b/services/core/internal/runtimeobs/otlpexporter/exporter.go index f4771a059..ddc395bd5 100644 --- a/services/core/internal/runtimeobs/otlpexporter/exporter.go +++ b/services/core/internal/runtimeobs/otlpexporter/exporter.go @@ -10,14 +10,14 @@ import ( "regexp" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetrichttp" "go.opentelemetry.io/otel/sdk/instrumentation" sdkmetric "go.opentelemetry.io/otel/sdk/metric" "go.opentelemetry.io/otel/sdk/metric/metricdata" "go.opentelemetry.io/otel/sdk/resource" - - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" ) const scopeName = "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs/otlpexporter" @@ -161,16 +161,16 @@ func validateRecord(record runtimeobs.ExportRecord) error { return errors.New("invalid Runtime history result reason") } switch record.Mode { - case runtimeobs.ModeNone, runtimeobs.ModeSelfHosted, runtimeobs.ModeManaged: + case v1.RuntimeModeNone, v1.RuntimeModeSelfHosted, v1.RuntimeModeManaged: default: return errors.New("invalid Runtime history mode") } switch record.Status { - case runtimeobs.StatusObserved: + case v1.RuntimeStatusObserved: if record.Sample == nil { return errors.New("observed Runtime history record has no sample") } - case runtimeobs.StatusUnavailable, runtimeobs.StatusUnsupported: + case v1.RuntimeStatusUnavailable, v1.RuntimeStatusUnsupported: if record.Sample != nil { return errors.New("unobserved Runtime history record has a sample") } diff --git a/services/core/internal/runtimeobs/otlpexporter/exporter_test.go b/services/core/internal/runtimeobs/otlpexporter/exporter_test.go index 0974de22c..389bd8926 100644 --- a/services/core/internal/runtimeobs/otlpexporter/exporter_test.go +++ b/services/core/internal/runtimeobs/otlpexporter/exporter_test.go @@ -9,12 +9,12 @@ import ( "testing" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/sdk/metric/metricdata" collectorv1 "go.opentelemetry.io/proto/otlp/collector/metrics/v1" "google.golang.org/protobuf/proto" - - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" ) type captureClient struct { @@ -41,7 +41,7 @@ func observedRecord() runtimeobs.ExportRecord { memoryUsage, memoryLimit := uint64(1024), uint64(2048) return runtimeobs.ExportRecord{ TenantID: "tenant", SessionID: "session", EnvironmentID: "environment", AllocationID: "allocation", - Mode: runtimeobs.ModeManaged, ProviderType: "docker", Status: runtimeobs.StatusObserved, + Mode: v1.RuntimeModeManaged, ProviderType: "docker", Status: v1.RuntimeStatusObserved, CollectionSource: runtimeobs.CollectionSourceOnRead, ResolvedAt: observedAt, SourceDuration: 50 * time.Millisecond, Sample: &runtimeobs.Sample{ @@ -110,7 +110,7 @@ func TestExporterPreservesUnavailableWithoutInventingResourceValues(t *testing.T exporter := newWithClient(client) record := runtimeobs.ExportRecord{ TenantID: "tenant", SessionID: "session", EnvironmentID: "environment", AllocationID: "allocation", - Mode: runtimeobs.ModeManaged, Status: runtimeobs.StatusUnavailable, Reason: "sample_timeout", + Mode: v1.RuntimeModeManaged, Status: v1.RuntimeStatusUnavailable, Reason: "sample_timeout", CollectionSource: runtimeobs.CollectionSourcePeriodic, ResolvedAt: time.Date(2026, 9, 23, 3, 0, 0, 0, time.UTC), } @@ -128,7 +128,7 @@ func TestExporterEmitsCanonicalSessionTokenGaugesWithoutProviderValues(t *testin exporter := newWithClient(client) record := runtimeobs.ExportRecord{ TenantID: "tenant", SessionID: "session", EnvironmentID: "environment", - Mode: runtimeobs.ModeManaged, Status: runtimeobs.StatusUnavailable, Reason: "sample_timeout", + Mode: v1.RuntimeModeManaged, Status: v1.RuntimeStatusUnavailable, Reason: "sample_timeout", CollectionSource: runtimeobs.CollectionSourcePeriodic, ResolvedAt: time.Date(2026, 9, 23, 3, 0, 0, 0, time.UTC), TokenUsage: &runtimeobs.TokenUsage{InputTokens: 120, OutputTokens: 30}, diff --git a/services/core/internal/runtimeobs/sample.go b/services/core/internal/runtimeobs/sample.go index b5e06271a..e1e318c44 100644 --- a/services/core/internal/runtimeobs/sample.go +++ b/services/core/internal/runtimeobs/sample.go @@ -8,22 +8,6 @@ import ( "time" ) -type Mode string - -const ( - ModeNone Mode = "none" - ModeSelfHosted Mode = "self_hosted" - ModeManaged Mode = "openai_hosted" -) - -type Status string - -const ( - StatusObserved Status = "observed" - StatusUnsupported Status = "unsupported" - StatusUnavailable Status = "unavailable" -) - // Sample contains provider-neutral cumulative counters and current gauges. // Pointer fields distinguish an observed zero from an unavailable measurement. type Sample struct { diff --git a/services/core/internal/runtimeobs/sampler_test.go b/services/core/internal/runtimeobs/sampler_test.go index da14b77b4..ae173f00d 100644 --- a/services/core/internal/runtimeobs/sampler_test.go +++ b/services/core/internal/runtimeobs/sampler_test.go @@ -7,6 +7,8 @@ import ( "sync/atomic" "testing" "time" + + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" ) type samplerLister struct { @@ -145,7 +147,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"}} + target := Target{EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} service, err := NewService(fixedResolver{target: target}, sourceOf(source)) if err != nil { t.Fatal(err) @@ -239,7 +241,7 @@ func TestSamplerPreservesProviderTimeoutAndFinalFenceAfterSlowResolution(t *test samplerLister: lister, delay: historyOwnershipCheckTimeout + 25*time.Millisecond, target: Target{ - EnvironmentID: "environment", Mode: ModeManaged, + EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}, }, } @@ -265,7 +267,7 @@ func TestSamplerPreservesProviderTimeoutAndFinalFenceAfterSlowResolution(t *test } select { case record := <-records: - if record.Status != StatusUnavailable || record.Reason != "sample_timeout" || record.CollectionSource != CollectionSourcePeriodic { + if record.Status != v1.RuntimeStatusUnavailable || record.Reason != "sample_timeout" || record.CollectionSource != CollectionSourcePeriodic { t.Fatalf("slow-resolution timeout export mismatch: %+v", record) } case <-time.After(time.Second): diff --git a/services/core/internal/runtimeobs/service.go b/services/core/internal/runtimeobs/service.go index 26e6c15a2..56dcab639 100644 --- a/services/core/internal/runtimeobs/service.go +++ b/services/core/internal/runtimeobs/service.go @@ -7,12 +7,13 @@ import ( "sync" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/providercontract" ) type Observation struct { Target Target - Status Status + Status v1.RuntimeObservationStatus Sample *Sample Reason string ProviderType string @@ -190,10 +191,10 @@ func (s *Service) resolve(ctx context.Context, tenantID, sessionID string, owner target, err := s.resolver.Resolve(ctx, tenantID, sessionID) resolvedAt := s.now() if errors.Is(err, ErrUnavailable) { - if target.TenantID != tenantID || target.SessionID != sessionID || target.Mode != ModeManaged || target.EnvironmentID == "" { + if target.TenantID != tenantID || target.SessionID != sessionID || target.Mode != v1.RuntimeModeManaged || target.EnvironmentID == "" { return Observation{}, nil, errors.New("Runtime observation resolver returned invalid pending allocation identity") } - return Observation{Target: target, Status: StatusUnavailable, Reason: "allocation_pending", ResolvedAt: resolvedAt}, nil, nil + return Observation{Target: target, Status: v1.RuntimeStatusUnavailable, Reason: "allocation_pending", ResolvedAt: resolvedAt}, nil, nil } if err != nil { return Observation{}, nil, err @@ -201,14 +202,14 @@ func (s *Service) resolve(ctx context.Context, tenantID, sessionID string, owner if target.TenantID != tenantID || target.SessionID != sessionID { return Observation{}, nil, errors.New("Runtime observation resolver returned mismatched ownership") } - if (target.Mode == ModeNone && target.EnvironmentID != "") || - ((target.Mode == ModeSelfHosted || target.Mode == ModeManaged) && target.EnvironmentID == "") { + if (target.Mode == v1.RuntimeModeNone && target.EnvironmentID != "") || + ((target.Mode == v1.RuntimeModeSelfHosted || target.Mode == v1.RuntimeModeManaged) && target.EnvironmentID == "") { return Observation{}, nil, errors.New("Runtime observation resolver returned mismatched Environment identity") } - if target.Mode == ModeNone || target.Mode == ModeSelfHosted { - return Observation{Target: target, Status: StatusUnsupported, Reason: "runtime_mode_not_observable", ResolvedAt: resolvedAt}, nil, nil + if target.Mode == v1.RuntimeModeNone || target.Mode == v1.RuntimeModeSelfHosted { + return Observation{Target: target, Status: v1.RuntimeStatusUnsupported, Reason: "runtime_mode_not_observable", ResolvedAt: resolvedAt}, nil, nil } - if target.Mode != ModeManaged || target.Instance.AllocationID == "" || target.Instance.ProviderKey == "" { + if target.Mode != v1.RuntimeModeManaged || target.Instance.AllocationID == "" || target.Instance.ProviderKey == "" { return Observation{}, nil, errors.New("invalid managed Runtime observation target") } if !target.Instance.AllocationCreatedAt.IsZero() && @@ -217,9 +218,9 @@ func (s *Service) resolve(ctx context.Context, tenantID, sessionID string, owner } switch target.Instance.AllocationState { case "creating": - return Observation{Target: target, Status: StatusUnavailable, Reason: "allocation_pending", ResolvedAt: resolvedAt}, nil, nil + return Observation{Target: target, Status: v1.RuntimeStatusUnavailable, Reason: "allocation_pending", ResolvedAt: resolvedAt}, nil, nil case "cleanup_pending", "released": - return Observation{Target: target, Status: StatusUnavailable, Reason: "runtime_not_running", ResolvedAt: resolvedAt}, nil, nil + return Observation{Target: target, Status: v1.RuntimeStatusUnavailable, Reason: "runtime_not_running", ResolvedAt: resolvedAt}, nil, nil case "running": default: return Observation{}, nil, errors.New("invalid managed Runtime allocation state") @@ -229,14 +230,14 @@ func (s *Service) resolve(ctx context.Context, tenantID, sessionID string, owner // complete classifies one provider result and hands it to history export. func (s *Service) complete(ctx context.Context, read *sourceRead, sample Sample, err error, sourceDuration time.Duration, collectionSource CollectionSource, owner OwnershipChecker) (Observation, error) { - observation := Observation{Target: read.target, Status: StatusUnavailable, ProviderType: read.providerType, SourceDuration: sourceDuration} + observation := Observation{Target: read.target, Status: v1.RuntimeStatusUnavailable, ProviderType: read.providerType, SourceDuration: sourceDuration} switch { case errors.Is(err, providercontract.ErrUnsupported): reason, valid := providercontract.UnsupportedReason(err, "Observe") if !valid { return Observation{}, providercontract.ErrContract } - observation.Status, observation.Reason = StatusUnsupported, reason + observation.Status, observation.Reason = v1.RuntimeStatusUnsupported, reason case errors.Is(err, context.DeadlineExceeded): observation.Reason = "sample_timeout" case errors.Is(err, ErrNotRunning): @@ -249,7 +250,7 @@ func (s *Service) complete(ctx context.Context, read *sourceRead, sample Sample, if err := sample.validate(s.now()); err != nil { return Observation{}, err } - observation.Status, observation.Sample = StatusObserved, &sample + observation.Status, observation.Sample = v1.RuntimeStatusObserved, &sample } observation.ResolvedAt = s.now() return s.finish(ctx, observation, collectionSource, owner) diff --git a/services/core/internal/runtimeobs/service_test.go b/services/core/internal/runtimeobs/service_test.go index 40f15551b..0c7288e51 100644 --- a/services/core/internal/runtimeobs/service_test.go +++ b/services/core/internal/runtimeobs/service_test.go @@ -8,6 +8,8 @@ import ( "sync" "testing" "time" + + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" ) type fixedResolver struct { @@ -109,10 +111,10 @@ func (e *gatedExporter) callCount() int { } func TestServiceDoesNotCallSourcesForUnsupportedModes(t *testing.T) { - for _, mode := range []Mode{ModeNone, ModeSelfHosted} { + for _, mode := range []v1.RuntimeObservationMode{v1.RuntimeModeNone, v1.RuntimeModeSelfHosted} { source := &fixedSource{} target := Target{Mode: mode} - if mode == ModeSelfHosted { + if mode == v1.RuntimeModeSelfHosted { target.EnvironmentID = "environment" } service, err := NewService(fixedResolver{target: target}, sourceOf(source)) @@ -120,14 +122,14 @@ func TestServiceDoesNotCallSourcesForUnsupportedModes(t *testing.T) { t.Fatal(err) } observation, err := service.ObserveSession(t.Context(), "tenant", "session") - if err != nil || observation.Status != StatusUnsupported || observation.Reason != "runtime_mode_not_observable" || source.calls != 0 { + if err != nil || observation.Status != v1.RuntimeStatusUnsupported || observation.Reason != "runtime_mode_not_observable" || source.calls != 0 { t.Fatalf("unsupported mode touched a source: %+v %v calls=%d", observation, err, source.calls) } } } func TestServicePreservesObservedZero(t *testing.T) { - target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} + target := Target{EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} zeroCPU := float64(0) zeroMemory := uint64(0) now := time.Date(2026, 9, 22, 1, 0, 0, 0, time.UTC) @@ -138,7 +140,7 @@ func TestServicePreservesObservedZero(t *testing.T) { } service.now = func() time.Time { return now } observation, err := service.ObserveSession(t.Context(), "tenant", "session") - if err != nil || observation.Status != StatusObserved || observation.Sample == nil || observation.Sample.CPUUsageSecondsTotal == nil || observation.Sample.MemoryUsageBytes == nil { + if err != nil || observation.Status != v1.RuntimeStatusObserved || observation.Sample == nil || observation.Sample.CPUUsageSecondsTotal == nil || observation.Sample.MemoryUsageBytes == nil { t.Fatalf("observed zero was lost: %+v %v", observation, err) } } @@ -151,7 +153,7 @@ func TestServiceExportsOnlySanitizedValidatedRecords(t *testing.T) { memoryUsage := uint64(1024) memoryLimit := uint64(2048) target := Target{ - TenantID: "tenant", SessionID: "session", EnvironmentID: "environment", Mode: ModeManaged, + TenantID: "tenant", SessionID: "session", EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{ AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running", ProviderState: json.RawMessage(`{"native_id":"must-not-export"}`), @@ -170,7 +172,7 @@ func TestServiceExportsOnlySanitizedValidatedRecords(t *testing.T) { service.now = func() time.Time { return now } observation, err := service.ObserveSession(t.Context(), "tenant", "session") - if err != nil || observation.Status != StatusObserved { + if err != nil || observation.Status != v1.RuntimeStatusObserved { t.Fatalf("observation failed: %+v %v", observation, err) } cpuSeconds = 99 @@ -181,7 +183,7 @@ func TestServiceExportsOnlySanitizedValidatedRecords(t *testing.T) { if record.TenantID != "tenant" || record.SessionID != "session" || record.EnvironmentID != "environment" || record.AllocationID != "allocation" { t.Fatalf("exported identity mismatch: %+v", record) } - if record.ProviderType != "docker" || record.Mode != ModeManaged || record.Status != StatusObserved || record.Reason != "" || record.CollectionSource != CollectionSourceOnRead { + if record.ProviderType != "docker" || record.Mode != v1.RuntimeModeManaged || record.Status != v1.RuntimeStatusObserved || record.Reason != "" || record.CollectionSource != CollectionSourceOnRead { t.Fatalf("exported classification mismatch: %+v", record) } if record.Sample == nil || record.Sample.CPUUsageSecondsTotal == nil || *record.Sample.CPUUsageSecondsTotal != 12.5 || record.Sample.MemoryUsageBytes == nil || *record.Sample.MemoryUsageBytes != 1024 { @@ -199,7 +201,7 @@ func TestServiceMarksPeriodicHistoryCollection(t *testing.T) { now := time.Date(2026, 9, 23, 4, 0, 0, 0, time.UTC) startedAt := now.Add(-time.Minute) target := Target{ - TenantID: "tenant", SessionID: "session", EnvironmentID: "environment", Mode: ModeManaged, + TenantID: "tenant", SessionID: "session", EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}, } records := make(chan ExportRecord, 1) @@ -232,7 +234,7 @@ func TestServiceDoesNotExportPeriodicSampleAfterOwnershipLoss(t *testing.T) { now := time.Date(2026, 9, 23, 4, 0, 0, 0, time.UTC) startedAt := now.Add(-time.Minute) target := Target{ - TenantID: "tenant", SessionID: "session", EnvironmentID: "environment", Mode: ModeManaged, + TenantID: "tenant", SessionID: "session", EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}, } records := make(chan ExportRecord, 1) @@ -261,7 +263,7 @@ func TestServiceDoesNotExportPeriodicSampleAfterOwnershipLoss(t *testing.T) { func TestServiceExportQueueNeverBlocksOrChangesObservation(t *testing.T) { now := time.Date(2026, 9, 22, 1, 0, 0, 0, time.UTC) - target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} + target := Target{EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} exporter := &gatedExporter{started: make(chan struct{}, 1), release: make(chan struct{})} service, err := NewService( fixedResolver{target: target}, @@ -273,7 +275,7 @@ func TestServiceExportQueueNeverBlocksOrChangesObservation(t *testing.T) { } service.now = func() time.Time { return now } - if observation, err := service.ObserveSession(t.Context(), "tenant", "session"); err != nil || observation.Status != StatusObserved { + if observation, err := service.ObserveSession(t.Context(), "tenant", "session"); err != nil || observation.Status != v1.RuntimeStatusObserved { t.Fatalf("first observation failed: %+v %v", observation, err) } select { @@ -285,7 +287,7 @@ func TestServiceExportQueueNeverBlocksOrChangesObservation(t *testing.T) { completed := make(chan error, 1) go func() { observation, observeErr := service.ObserveSession(t.Context(), "tenant", "session") - if observeErr == nil && observation.Status != StatusObserved { + if observeErr == nil && observation.Status != v1.RuntimeStatusObserved { observeErr = errors.New("unexpected observation status") } completed <- observeErr @@ -324,7 +326,7 @@ func TestServiceIgnoresExporterFailureAndValidatesOptions(t *testing.T) { now := time.Date(2026, 9, 22, 1, 0, 0, 0, time.UTC) records := make(chan ExportRecord, 1) - target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} + target := Target{EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} service, err := NewService( fixedResolver{target: target}, sourceOf(&fixedSource{sample: Sample{ObservedAt: now}}), @@ -335,7 +337,7 @@ func TestServiceIgnoresExporterFailureAndValidatesOptions(t *testing.T) { } service.now = func() time.Time { return now } observation, err := service.ObserveSession(t.Context(), "tenant", "session") - if err != nil || observation.Status != StatusObserved { + if err != nil || observation.Status != v1.RuntimeStatusObserved { t.Fatalf("exporter failure changed observation: %+v %v", observation, err) } if err := service.Close(t.Context()); err != nil { @@ -346,7 +348,7 @@ func TestServiceIgnoresExporterFailureAndValidatesOptions(t *testing.T) { func TestServiceCloseHonorsItsDeadlineWhenExporterDoesNot(t *testing.T) { now := time.Date(2026, 9, 22, 1, 0, 0, 0, time.UTC) exporter := stubbornExporter{started: make(chan struct{}), release: make(chan struct{})} - target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} + target := Target{EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} service, err := NewService( fixedResolver{target: target}, sourceOf(&fixedSource{sample: Sample{ObservedAt: now}}), @@ -376,7 +378,7 @@ func TestServiceCloseHonorsItsDeadlineWhenExporterDoesNot(t *testing.T) { func TestServiceIsolatesExporterPanics(t *testing.T) { now := time.Date(2026, 9, 22, 1, 0, 0, 0, time.UTC) exporter := panicExporter{called: make(chan struct{}, 1)} - target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} + target := Target{EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} service, err := NewService( fixedResolver{target: target}, sourceOf(&fixedSource{sample: Sample{ObservedAt: now}}), @@ -387,7 +389,7 @@ func TestServiceIsolatesExporterPanics(t *testing.T) { } service.now = func() time.Time { return now } observation, err := service.ObserveSession(t.Context(), "tenant", "session") - if err != nil || observation.Status != StatusObserved { + if err != nil || observation.Status != v1.RuntimeStatusObserved { t.Fatalf("observation failed: %+v %v", observation, err) } if err := service.Close(t.Context()); err != nil { @@ -401,7 +403,7 @@ func TestServiceIsolatesExporterPanics(t *testing.T) { } func TestServiceMapsOnlyDeclaredUnavailability(t *testing.T) { - target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} + target := Target{EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} for _, tc := range []struct { err error wantReason string @@ -420,7 +422,7 @@ func TestServiceMapsOnlyDeclaredUnavailability(t *testing.T) { if (err != nil) != tc.wantError { t.Fatalf("wrong error classification: %+v %v", observation, err) } - if !tc.wantError && (observation.Status != StatusUnavailable || observation.Reason != tc.wantReason || observation.ProviderType != "docker") { + if !tc.wantError && (observation.Status != v1.RuntimeStatusUnavailable || observation.Reason != tc.wantReason || observation.ProviderType != "docker") { t.Fatalf("declared unavailability was not mapped: %+v", observation) } } @@ -428,7 +430,7 @@ func TestServiceMapsOnlyDeclaredUnavailability(t *testing.T) { func TestServiceMapsAnActualSourceDeadlineWithoutLeakingIt(t *testing.T) { target := Target{ - EnvironmentID: "environment", Mode: ModeManaged, + EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}, } service, err := NewService(fixedResolver{target: target}, sourceOf(blockingSource{})) @@ -438,14 +440,14 @@ func TestServiceMapsAnActualSourceDeadlineWithoutLeakingIt(t *testing.T) { ctx, cancel := context.WithTimeout(t.Context(), 10*time.Millisecond) defer cancel() observation, err := service.ObserveSession(ctx, "tenant", "session") - if err != nil || observation.Status != StatusUnavailable || observation.Reason != "sample_timeout" || observation.ProviderType != "docker" { + if err != nil || observation.Status != v1.RuntimeStatusUnavailable || observation.Reason != "sample_timeout" || observation.ProviderType != "docker" { t.Fatalf("source deadline was not safely classified: %+v %v", observation, err) } } func TestServiceExportsPeriodicSourceTimeoutAfterFinalOwnershipFence(t *testing.T) { target := Target{ - TenantID: "tenant", SessionID: "session", EnvironmentID: "environment", Mode: ModeManaged, + TenantID: "tenant", SessionID: "session", EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}, } records := make(chan ExportRecord, 1) @@ -459,7 +461,7 @@ func TestServiceExportsPeriodicSourceTimeoutAfterFinalOwnershipFence(t *testing. } owner := &sequenceOwner{} observation, err := service.ObserveSessionForHistory(t.Context(), "tenant", "session", owner, 10*time.Millisecond) - if err != nil || observation.Status != StatusUnavailable || observation.Reason != "sample_timeout" { + if err != nil || observation.Status != v1.RuntimeStatusUnavailable || observation.Reason != "sample_timeout" { t.Fatalf("periodic source timeout was not safely classified: %+v %v", observation, err) } if calls := owner.calls.Load(); calls != 2 { @@ -467,7 +469,7 @@ func TestServiceExportsPeriodicSourceTimeoutAfterFinalOwnershipFence(t *testing. } select { case record := <-records: - if record.Status != StatusUnavailable || record.Reason != "sample_timeout" || record.CollectionSource != CollectionSourcePeriodic { + if record.Status != v1.RuntimeStatusUnavailable || record.Reason != "sample_timeout" || record.CollectionSource != CollectionSourcePeriodic { t.Fatalf("periodic timeout export mismatch: %+v", record) } case <-time.After(time.Second): @@ -480,39 +482,39 @@ func TestServiceExportsPeriodicSourceTimeoutAfterFinalOwnershipFence(t *testing. func TestServiceClassifiesResolverAndTerminalAllocationUnavailability(t *testing.T) { now := time.Date(2026, 9, 22, 1, 0, 0, 0, time.UTC) - service, err := NewService(fixedResolver{target: Target{SessionID: "session", EnvironmentID: "environment", Mode: ModeManaged}, err: ErrUnavailable}, sourceOf(&fixedSource{})) + service, err := NewService(fixedResolver{target: Target{SessionID: "session", EnvironmentID: "environment", Mode: v1.RuntimeModeManaged}, err: ErrUnavailable}, sourceOf(&fixedSource{})) if err != nil { t.Fatal(err) } service.now = func() time.Time { return now } observation, err := service.ObserveSession(t.Context(), "tenant", "session") - if err != nil || observation.Status != StatusUnavailable || observation.Reason != "allocation_pending" || !observation.ResolvedAt.Equal(now) { + if err != nil || observation.Status != v1.RuntimeStatusUnavailable || observation.Reason != "allocation_pending" || !observation.ResolvedAt.Equal(now) { t.Fatalf("pending allocation was not classified: %+v %v", observation, err) } source := &fixedSource{} - target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "creating"}} + target := Target{EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "creating"}} service, err = NewService(fixedResolver{target: target}, sourceOf(source)) if err != nil { t.Fatal(err) } observation, err = service.ObserveSession(t.Context(), "tenant", "session") - if err != nil || observation.Status != StatusUnavailable || observation.Reason != "allocation_pending" || source.calls != 0 { + if err != nil || observation.Status != v1.RuntimeStatusUnavailable || observation.Reason != "allocation_pending" || source.calls != 0 { t.Fatalf("creating allocation reached its provider: %+v %v calls=%d", observation, err, source.calls) } - target = Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "released"}} + target = Target{EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "released"}} service, err = NewService(fixedResolver{target: target}, sourceOf(&fixedSource{})) if err != nil { t.Fatal(err) } observation, err = service.ObserveSession(t.Context(), "tenant", "session") - if err != nil || observation.Status != StatusUnavailable || observation.Reason != "runtime_not_running" { + if err != nil || observation.Status != v1.RuntimeStatusUnavailable || observation.Reason != "runtime_not_running" { t.Fatalf("released allocation was not classified: %+v %v", observation, err) } source = &fixedSource{} - target = Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{ + target = Target{EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{ AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running", AllocationCreatedAt: time.Now().Add(time.Hour), }} service, err = NewService(fixedResolver{target: target}, sourceOf(source)) @@ -525,7 +527,7 @@ func TestServiceClassifiesResolverAndTerminalAllocationUnavailability(t *testing } func TestServiceRejectsUnsafeProviderSamples(t *testing.T) { - target := Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} + target := Target{EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}} now := time.Date(2026, 9, 22, 1, 0, 0, 0, time.UTC) preEpoch := time.Unix(-1, 0).UTC() tooLarge := uint64(1 << 53) @@ -552,10 +554,10 @@ func TestServiceRejectsUnsafeProviderSamples(t *testing.T) { func TestServiceRejectsMismatchedResolvedOwnership(t *testing.T) { for _, target := range []Target{ - {TenantID: "other", SessionID: "session", Mode: ModeNone}, - {TenantID: "tenant", SessionID: "other", Mode: ModeNone}, - {TenantID: "tenant", SessionID: "session", EnvironmentID: "unexpected", Mode: ModeNone}, - {TenantID: "tenant", SessionID: "session", Mode: ModeSelfHosted}, + {TenantID: "other", SessionID: "session", Mode: v1.RuntimeModeNone}, + {TenantID: "tenant", SessionID: "other", Mode: v1.RuntimeModeNone}, + {TenantID: "tenant", SessionID: "session", EnvironmentID: "unexpected", Mode: v1.RuntimeModeNone}, + {TenantID: "tenant", SessionID: "session", Mode: v1.RuntimeModeSelfHosted}, } { service, err := NewService(fixedResolver{target: target}, sourceOf(&fixedSource{})) if err != nil { diff --git a/services/core/internal/runtimeobs/source_test.go b/services/core/internal/runtimeobs/source_test.go index 655d5a64e..d87d8a41a 100644 --- a/services/core/internal/runtimeobs/source_test.go +++ b/services/core/internal/runtimeobs/source_test.go @@ -6,6 +6,7 @@ import ( "testing" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/providercontract" ) @@ -22,7 +23,7 @@ func (s *selectingSource) load(context.Context) (Source, string, error) { } func observableTarget() fixedResolver { - return fixedResolver{target: Target{EnvironmentID: "environment", Mode: ModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}}} + return fixedResolver{target: Target{EnvironmentID: "environment", Mode: v1.RuntimeModeManaged, Instance: Instance{AllocationID: "allocation", ProviderKey: "provider", AllocationState: "running"}}} } func TestUnavailableAndNilSourceSelection(t *testing.T) { @@ -32,7 +33,7 @@ func TestUnavailableAndNilSourceSelection(t *testing.T) { 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 != "" { + if err != nil || observation.Status != v1.RuntimeStatusUnavailable || observation.Reason != "sample_unavailable" || observation.ProviderType != "" { t.Fatal(observation, err) } resolver.err = nil @@ -68,7 +69,7 @@ func TestSourceSelectionStaysBoundForWholePage(t *testing.T) { } 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 { + if errs[index] != nil || observation.ProviderType != "previous" || observation.Status != v1.RuntimeStatusObserved { t.Fatal(index, observation, errs[index]) } } diff --git a/services/core/internal/sandbox/docker/provider_test.go b/services/core/internal/sandbox/docker/provider_test.go index 73a83dda0..56bfa269b 100644 --- a/services/core/internal/sandbox/docker/provider_test.go +++ b/services/core/internal/sandbox/docker/provider_test.go @@ -10,6 +10,7 @@ import ( "testing" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxbootstrap" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" @@ -77,7 +78,7 @@ func TestDockerProviderLifecycle(t *testing.T) { info, e := p.Create(ctx, b) contracttest.AssertObservation(t, info, e, b.Reference, "", "running") resources, e := p.Observe(ctx, runtimeobs.Target{ - TenantID: b.TenantID, EnvironmentID: b.EnvironmentID, Mode: runtimeobs.ModeManaged, + TenantID: b.TenantID, EnvironmentID: b.EnvironmentID, Mode: v1.RuntimeModeManaged, Instance: runtimeobs.Instance{AllocationID: b.AllocationID, ProviderKey: installationID}, }) if e != nil || resources.StartedAt == nil || resources.CPUUsageSecondsTotal == nil || resources.MemoryUsageBytes == nil || resources.CPUCapacityCores == nil || resources.MemoryLimitBytes == nil { diff --git a/services/core/internal/sandbox/docker/resources.go b/services/core/internal/sandbox/docker/resources.go index 863af750a..b9e141f72 100644 --- a/services/core/internal/sandbox/docker/resources.go +++ b/services/core/internal/sandbox/docker/resources.go @@ -7,6 +7,7 @@ import ( "fmt" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" "github.com/moby/moby/api/types/container" @@ -16,7 +17,7 @@ import ( // Observe is read-only. Inspect verifies allocation ownership before Docker // statistics are requested; it never renews or changes the container. func (p *Provider) Observe(ctx context.Context, target runtimeobs.Target) (runtimeobs.Sample, error) { - if target.Mode != runtimeobs.ModeManaged || target.Instance.AllocationID == "" { + if target.Mode != v1.RuntimeModeManaged || target.Instance.AllocationID == "" { return runtimeobs.Sample{}, sandbox.ErrInvalid } if target.Instance.ProviderKey != p.config.InstallationID { diff --git a/services/core/internal/sandbox/docker/resources_test.go b/services/core/internal/sandbox/docker/resources_test.go index 4911716d4..c6995aa85 100644 --- a/services/core/internal/sandbox/docker/resources_test.go +++ b/services/core/internal/sandbox/docker/resources_test.go @@ -9,6 +9,7 @@ import ( "testing" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" "github.com/google/uuid" "github.com/moby/moby/api/types/container" @@ -18,7 +19,7 @@ import ( func TestObserveVerifiesOwnershipThenReadsOneShotStats(t *testing.T) { installationID := uuid.NewString() target := runtimeobs.Target{ - TenantID: uuid.NewString(), SessionID: uuid.NewString(), EnvironmentID: uuid.NewString(), Mode: runtimeobs.ModeManaged, + TenantID: uuid.NewString(), SessionID: uuid.NewString(), EnvironmentID: uuid.NewString(), Mode: v1.RuntimeModeManaged, Instance: runtimeobs.Instance{AllocationID: uuid.NewString(), ProviderKey: installationID}, } observed := time.Now().UTC().Truncate(time.Microsecond) diff --git a/services/core/internal/sandbox/e2b/observations.go b/services/core/internal/sandbox/e2b/observations.go index 4ea46e064..95fda3482 100644 --- a/services/core/internal/sandbox/e2b/observations.go +++ b/services/core/internal/sandbox/e2b/observations.go @@ -5,6 +5,7 @@ import ( "math" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" ) @@ -32,7 +33,7 @@ type Observation struct { // changes a sandbox. func (p *Provider) Observe(ctx context.Context, target runtimeobs.Target) (runtimeobs.Sample, error) { reference := sandbox.Reference{TenantID: target.TenantID, EnvironmentID: target.EnvironmentID, AllocationID: target.Instance.AllocationID} - if target.Mode != runtimeobs.ModeManaged || !validReference(reference) { + if target.Mode != v1.RuntimeModeManaged || !validReference(reference) { return runtimeobs.Sample{}, sandbox.ErrInvalid } if target.Instance.ProviderKey != p.config.InstallationID { diff --git a/services/core/internal/sandbox/e2b/observations_test.go b/services/core/internal/sandbox/e2b/observations_test.go index 04acefd4d..10e104e3a 100644 --- a/services/core/internal/sandbox/e2b/observations_test.go +++ b/services/core/internal/sandbox/e2b/observations_test.go @@ -5,6 +5,7 @@ import ( "testing" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" ) @@ -16,7 +17,7 @@ func TestObserveMapsMetricsAndKeepsUnmeasuredValuesNull(t *testing.T) { count, percent, memoryUsed, memoryTotal, diskUsed, diskTotal := 2.0, 19.55, uint64(183836672), uint64(2079141888), uint64(1593188352), uint64(23511863296) caller.response.Observation = &Observation{Status: "observed", ObservedAt: &observed, StartedAt: &started, CPUCount: &count, CPUUsedPct: &percent, MemUsed: &memoryUsed, MemTotal: &memoryTotal, DiskUsed: &diskUsed, DiskTotal: &diskTotal} - target := runtimeobs.Target{TenantID: running.TenantID, EnvironmentID: running.EnvironmentID, Mode: runtimeobs.ModeManaged, + target := runtimeobs.Target{TenantID: running.TenantID, EnvironmentID: running.EnvironmentID, Mode: v1.RuntimeModeManaged, Instance: runtimeobs.Instance{AllocationID: running.AllocationID, ProviderKey: p.config.InstallationID}} sample, err := p.Observe(bounded(t), target) if err != nil || len(caller.requests) != 1 || caller.requests[0].Operation != "observe" || diff --git a/services/core/internal/sandbox/microsandbox/resources.go b/services/core/internal/sandbox/microsandbox/resources.go index ce586b389..e5d0e38bc 100644 --- a/services/core/internal/sandbox/microsandbox/resources.go +++ b/services/core/internal/sandbox/microsandbox/resources.go @@ -6,6 +6,7 @@ import ( "errors" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" ) @@ -14,7 +15,7 @@ import ( // one-shot helper. The persisted compute receipt selects the exact generation; // browser input and provider display names never select a sandbox. func (p *Provider) Observe(ctx context.Context, target runtimeobs.Target) (runtimeobs.Sample, error) { - if target.Mode != runtimeobs.ModeManaged || target.Instance.AllocationID == "" { + if target.Mode != v1.RuntimeModeManaged || target.Instance.AllocationID == "" { return runtimeobs.Sample{}, sandbox.ErrInvalid } if target.Instance.ProviderKey != p.config.InstallationID { diff --git a/services/core/internal/sandbox/microsandbox/resources_test.go b/services/core/internal/sandbox/microsandbox/resources_test.go index c52eaf2ec..245f57c15 100644 --- a/services/core/internal/sandbox/microsandbox/resources_test.go +++ b/services/core/internal/sandbox/microsandbox/resources_test.go @@ -7,6 +7,7 @@ import ( "testing" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" ) @@ -21,7 +22,7 @@ func observationTarget(t *testing.T, compute Compute) runtimeobs.Target { } r := testRef() return runtimeobs.Target{ - TenantID: r.TenantID, SessionID: "55555555-5555-4555-8555-555555555555", EnvironmentID: r.EnvironmentID, Mode: runtimeobs.ModeManaged, + TenantID: r.TenantID, SessionID: "55555555-5555-4555-8555-555555555555", EnvironmentID: r.EnvironmentID, Mode: v1.RuntimeModeManaged, Instance: runtimeobs.Instance{AllocationID: r.AllocationID, ProviderKey: testConfig().InstallationID, AllocationState: "running", ComputePhase: "running", ProviderState: state}, } } diff --git a/services/core/internal/sandbox/node/observations.go b/services/core/internal/sandbox/node/observations.go index 035869f46..40f458500 100644 --- a/services/core/internal/sandbox/node/observations.go +++ b/services/core/internal/sandbox/node/observations.go @@ -4,6 +4,7 @@ import ( "context" "errors" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" ) @@ -13,7 +14,7 @@ func observationReference(target runtimeobs.Target) sandbox.Reference { } func validObservation(target runtimeobs.Target, reference sandbox.Reference) bool { - return target.Mode == runtimeobs.ModeManaged && validID(target.SessionID) && + return target.Mode == v1.RuntimeModeManaged && validID(target.SessionID) && validID(target.Instance.ProviderKey) && observationReference(target) == reference && target.TokenUsage == nil } diff --git a/services/core/internal/sandbox/node/observations_test.go b/services/core/internal/sandbox/node/observations_test.go index c5ca75325..9392ca001 100644 --- a/services/core/internal/sandbox/node/observations_test.go +++ b/services/core/internal/sandbox/node/observations_test.go @@ -10,11 +10,11 @@ import ( "testing" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/providercontract" - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox/docker" - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox/docker" "github.com/google/uuid" ) @@ -36,7 +36,7 @@ func (p *observationProvider) Observe(_ context.Context, target runtimeobs.Targe return p.sample, p.err } func observationTarget(r sandbox.Reference, installation string) runtimeobs.Target { - return runtimeobs.Target{TenantID: r.TenantID, SessionID: uuid.NewString(), EnvironmentID: r.EnvironmentID, Mode: runtimeobs.ModeManaged, + return runtimeobs.Target{TenantID: r.TenantID, SessionID: uuid.NewString(), EnvironmentID: r.EnvironmentID, Mode: v1.RuntimeModeManaged, Instance: runtimeobs.Instance{AllocationID: r.AllocationID, ProviderKey: installation, AllocationState: "running", ComputePhase: "running", ProviderState: json.RawMessage(`{"current":{"id":"exact-incarnation","generation":3}}`)}, TokenUsage: &runtimeobs.TokenUsage{InputTokens: 123, OutputTokens: 456}} } @@ -157,7 +157,7 @@ func TestObservationWirePreservesUnavailableAndRejectsMismatchedIdentity(t *test }) } p := &observationProvider{fakeProvider: &fakeProvider{}} - for _, mutate := range []func(*runtimeobs.Target){func(t *runtimeobs.Target) { t.EnvironmentID = uuid.NewString() }, func(t *runtimeobs.Target) { t.Instance.AllocationID = uuid.NewString() }, func(t *runtimeobs.Target) { t.TenantID = uuid.NewString() }, func(t *runtimeobs.Target) { t.TokenUsage = &runtimeobs.TokenUsage{} }, func(t *runtimeobs.Target) { t.Mode = runtimeobs.ModeSelfHosted }} { + for _, mutate := range []func(*runtimeobs.Target){func(t *runtimeobs.Target) { t.EnvironmentID = uuid.NewString() }, func(t *runtimeobs.Target) { t.Instance.AllocationID = uuid.NewString() }, func(t *runtimeobs.Target) { t.TenantID = uuid.NewString() }, func(t *runtimeobs.Target) { t.TokenUsage = &runtimeobs.TokenUsage{} }, func(t *runtimeobs.Target) { t.Mode = v1.RuntimeModeSelfHosted }} { invalid := target mutate(&invalid) q.Observation = &invalid From 9df86f5a8bd0b6eab96b2f4e3afc82f8f4107b57 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Thu, 8 Oct 2026 15:11:43 +0000 Subject: [PATCH 2/2] Sync microsandbox helper observation dependencies --- services/core/tools/microsandbox-provider/go.mod | 4 ++++ services/core/tools/microsandbox-provider/go.sum | 6 ++++++ 2 files changed, 10 insertions(+) diff --git a/services/core/tools/microsandbox-provider/go.mod b/services/core/tools/microsandbox-provider/go.mod index 1483c8963..560f15ac4 100644 --- a/services/core/tools/microsandbox-provider/go.mod +++ b/services/core/tools/microsandbox-provider/go.mod @@ -12,7 +12,11 @@ require ( github.com/gorilla/websocket v1.5.3 // indirect github.com/libp2p/go-buffer-pool v0.0.2 // indirect github.com/libp2p/go-yamux/v5 v5.1.0 // indirect + golang.org/x/net v0.58.0 // indirect golang.org/x/sync v0.22.0 // indirect + golang.org/x/sys v0.47.0 // indirect + golang.org/x/text v0.41.0 // indirect + gopkg.in/yaml.v3 v3.0.1 // indirect ) // The helper ships from this repository; the SDK itself is a released module. diff --git a/services/core/tools/microsandbox-provider/go.sum b/services/core/tools/microsandbox-provider/go.sum index 28c8c14f3..e00757722 100644 --- a/services/core/tools/microsandbox-provider/go.sum +++ b/services/core/tools/microsandbox-provider/go.sum @@ -54,9 +54,15 @@ go.opentelemetry.io/otel/metric v1.44.0 h1:1w0gILTcHdr3YI+ixLyjemwrVnsMURbTZFrSY go.opentelemetry.io/otel/metric v1.44.0/go.mod h1:8O7hanEPBNgEMmybD3s2VBKcgWOCsA6tzHBPODAiquo= go.opentelemetry.io/otel/trace v1.44.0 h1:jxF5CsGYCe74MCRx2X4g7WsY/VBKRqqpNvXlX/6gtIk= go.opentelemetry.io/otel/trace v1.44.0/go.mod h1:oLl1jrMQAVo6v3GAggN+1VH9VIz9iUSvW53sW1Q8PIE= +golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To= +golang.org/x/net v0.58.0/go.mod h1:YwCddHnFlT7eLQqVprV19OnhLGtc5xOKgE0RyqgfWAU= golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek= golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs= golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/text v0.41.0 h1:vz/seA0lnX87Othu2f/0L24RcgrXD9/YFTSuGjj3rH8= +golang.org/x/text v0.41.0/go.mod h1:jvf1O8ajNzZqhSrQBPbutR/EB83Cc0CFrezNQIwbb5M= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=