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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
The table of contents is too big for display.
Diff view
Diff view
  •  
  •  
  •  
32 changes: 16 additions & 16 deletions contracts/agents-api/core.openapi.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -196,7 +196,7 @@ definitions:
agent_id:
type: string
assets:
$ref: '#/definitions/store.AdminAssetCounts'
$ref: '#/definitions/sessions.AdminAssetCounts'
coverage:
$ref: '#/definitions/api.AdminUsageCoverage'
key_id:
Expand Down Expand Up @@ -1395,6 +1395,21 @@ definitions:
source_commit:
type: string
type: object
sessions.AdminAssetCounts:
properties:
agents:
type: integer
credentials:
type: integer
environment_templates:
type: integer
files:
type: integer
skills:
type: integer
vaults:
type: integer
type: object
sessions.ExecutorCredential:
properties:
created_at:
Expand Down Expand Up @@ -1424,21 +1439,6 @@ definitions:
state:
type: string
type: object
store.AdminAssetCounts:
properties:
agents:
type: integer
credentials:
type: integer
environment_templates:
type: integer
files:
type: integer
skills:
type: integer
vaults:
type: integer
type: object
v1.Agent:
properties:
id:
Expand Down
2 changes: 1 addition & 1 deletion contracts/agents-api/runtime-observability.md
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ self-hosted: tenant_id -> session_id -> environment_id -> device_id + connection
none: tenant_id -> session_id (no Session-owned Runtime instance)
```

The resolver (`internal/runtimeobs/storeresolver`) reads the Session, its Environment, the current allocation and the Session's measured usage from the store. A Session, daemon connection, process, container and native Harness Session are different identities and never stand in for one another.
The resolver (`services/core/internal/deployment/observation.go`) reads the Session, its Environment, the current allocation and the Session's measured usage from the database. A Session, daemon connection, process, container and native Harness Session are different identities and never stand in for one another.

Managed Docker, microsandbox and E2B allocations are observed. `none` and `self_hosted` Sessions are `unsupported`; Core never attributes shared host statistics to an `environment:none` Session.

Expand Down
4 changes: 2 additions & 2 deletions contracts/agents-api/zh/runtime-observability.md
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
---
title: "运行时可观测性"
source: contracts/agents-api/runtime-observability.md
source_hash: 5da1279a81a6dcb3a85661dbba942c8937f871ae351ed80550db65b2a59156db
source_hash: f79a077350118f5f9bb3f4e8874242af3cb716e92f5e090f6652b001ae5780a1
---

这是面向贡献者的契约,规定 Core 如何观测 Runtime 并保留其历史。路由和响应字段见 [Runtime telemetry API](runtime-observability-api.md)。代码位于 `services/core/internal/runtimeobs`(解析、源、采样器和导出)、`internal/runtimehistory`(历史查询和 PostgreSQL 存储)以及 `internal/runtimeobs/otlpexporter`。
Expand All @@ -18,7 +18,7 @@ self-hosted: tenant_id -> session_id -> environment_id -> device_id + connection
none: tenant_id -> session_id (no Session-owned Runtime instance)
```

解析器(`internal/runtimeobs/storeresolver`)从存储中读取 Session、其 Environment、当前分配以及 Session 的实测使用量。Session、守护进程连接、进程、容器和原生 Harness Session 是不同身份,彼此绝不能替代。
解析器(`services/core/internal/deployment/observation.go`)从数据库中读取 Session、其 Environment、当前分配以及 Session 的实测使用量。Session、守护进程连接、进程、容器和原生 Harness Session 是不同身份,彼此绝不能替代。

托管 Docker、microsandbox 和 E2B 分配均会被观测。`none` 和 `self_hosted` Session 为 `unsupported`;Core 绝不会将共享主机统计信息归属于 `environment:none` Session。

Expand Down
14 changes: 8 additions & 6 deletions services/core/IMPLEMENTATION.md

Large diffs are not rendered by default.

18 changes: 5 additions & 13 deletions services/core/cmd/server/core_metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,8 @@ import (

"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/coremetrics"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/coremetricspg"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store"
"github.com/jackc/pgx/v5/pgxpool"
)

Expand All @@ -16,7 +16,7 @@ var buildRevision string
var processStartedAt = time.Now().UTC()

type coreMetricsSource struct {
store *store.Store
store *coremetricspg.Store
pool *pgxpool.Pool
worker *execution.Worker
registry *runtimegateway.Registry
Expand Down Expand Up @@ -51,7 +51,7 @@ func (s *coreMetricsSource) Sample(ctx context.Context) coremetrics.Sample {
if s.registry != nil {
devices = s.registry.Devices()
}
counts, err := s.store.ReadCoreExecutionSnapshot(ctx, time.Now(), devices)
counts, err := s.store.ReadExecutionSnapshot(ctx, time.Now(), devices)
if err != nil {
sample.Healthy = false
} else {
Expand All @@ -61,7 +61,7 @@ func (s *coreMetricsSource) Sample(ctx context.Context) coremetrics.Sample {
sample.WaitingForDaemon = metricPtr(counts.WaitingForDaemon)
}
}
size, err := s.store.ReadCoreDatabaseSize(ctx)
size, err := s.store.ReadDatabaseSize(ctx)
if err != nil {
sample.Healthy = false
} else {
Expand All @@ -71,15 +71,7 @@ func (s *coreMetricsSource) Sample(ctx context.Context) coremetrics.Sample {
return sample
}
func (s *coreMetricsSource) History(ctx context.Context, start, end time.Time, step time.Duration) (coremetrics.History, error) {
value, err := s.store.ReadCoreExecutionHistory(ctx, start, end, step)
if err != nil {
return coremetrics.History{}, err
}
result := coremetrics.History{Interrupted: value.Interrupted, QueueWaitMS: coremetrics.Latency{P50: value.QueueWaitMS.P50, P95: value.QueueWaitMS.P95}, Buckets: map[time.Time]*float64{}}
for _, bucket := range value.Buckets {
result.Buckets[bucket.Start.UTC()] = bucket.P95MS
}
return result, nil
return s.store.ReadExecutionHistory(ctx, start, end, step)
}

func reportCleanupResult(metrics *coremetrics.Service, job string, count int64, err error) {
Expand Down
1 change: 1 addition & 0 deletions services/core/cmd/server/http_routes_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,7 @@ func daemonComposition(t testing.TB) http.Handler {
Agents: struct{ api.Agents }{}, AgentsReader: struct{ api.AgentsReader }{},
EnvironmentTemplates: struct{ api.EnvironmentTemplates }{}, EnvironmentTemplatesReader: struct{ api.EnvironmentTemplatesReader }{},
Sessions: struct{ api.Sessions }{},
SessionsReader: struct{ api.SessionsReader }{},
SessionCreation: struct{ api.SessionCreation }{},
SessionEvents: struct{ api.SessionEvents }{},
Turns: struct{ api.Turns }{},
Expand Down
33 changes: 15 additions & 18 deletions services/core/cmd/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@ import (
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/nativeinstaller"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/agentpg"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/auditpg"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/coremetricspg"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/deploymentpg"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/filepg"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/modelconfigurationpg"
Expand All @@ -60,9 +61,7 @@ import (
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeenrollment"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimehistory"
historystoreresolver "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimehistory/storeresolver"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs"
observationstoreresolver "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs/storeresolver"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox/providers"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/skills"
Expand Down Expand Up @@ -198,7 +197,7 @@ func run() error {
if err != nil {
return err
}
metricsSource := &coreMetricsSource{store: executionStore, pool: pool}
metricsSource := &coreMetricsSource{store: coremetricspg.New(units), pool: pool}
metrics := coremetrics.New(processStartedAt, buildRevision, metricsSource)
auditRetention, err := writeAuditRetention()
if err != nil {
Expand Down Expand Up @@ -229,11 +228,11 @@ func run() error {
managed = managedNodes.runtime
observationSources[managed.InstallationID] = managedNodes.setup
}
observationResolver, err := observationstoreresolver.NewResolver(executionStore, deploymentStore)
observationResolver, err := deployment.NewObservationResolver(sessionStore, deploymentStore)
if err != nil {
return err
}
history, err := runtimeHistory(ctx, executionStore, public != "")
history, err := runtimeHistory(ctx, units, public != "")
if err != nil {
return err
}
Expand Down Expand Up @@ -273,11 +272,7 @@ func run() error {
if err := api.ValidateCredentialSeparation(ctx, keyAdmin, projectStore); err != nil {
return err
}
historyResolver, err := historystoreresolver.NewResolver(sessionStore)
if err != nil {
return err
}
historyService, err := runtimehistory.NewService(historyResolver, history.Reader)
historyService, err := runtimehistory.NewService(sessionStore, history.Reader)
if err != nil {
return err
}
Expand All @@ -290,7 +285,7 @@ func run() error {
if err != nil {
return err
}
daemonHandler, registry, err = runtime.NewGateway(sessionStore, sessionService, executionStore, executorURL)
daemonHandler, registry, err = runtime.NewGateway(sessionStore, sessionService, sessionStore, executorURL)
if err != nil {
return err
}
Expand All @@ -306,6 +301,7 @@ func run() error {
nativeInstaller = &api.NativeInstaller{Version: buildRevision, Catalog: catalog}
}
}
var deploymentExecution *deployment.ExecutionOperations
if registry != nil {
dispatcher := &execution.Dispatcher{Store: executionStore, Registry: registry,
Credentials: vaultService, Observer: modelConfigurationStore, Deployment: deploymentService, DeploymentReader: deploymentStore,
Expand All @@ -316,7 +312,7 @@ func run() error {
if err != nil {
return err
}
deploymentExecution, err := deployment.NewExecutionOperations(deploymentService, deploymentpg.NewExecution(lease, credentialKey))
deploymentExecution, err = deployment.NewExecutionOperations(deploymentService, deploymentpg.NewExecution(lease, credentialKey))
if err != nil {
return errors.Join(err, lease.Close(ctx))
}
Expand Down Expand Up @@ -402,25 +398,26 @@ func run() error {
EnvironmentTemplates: environmentTemplates, EnvironmentTemplatesReader: templateStore,
Files: fileService, FilesReader: fileStore,
Agents: agentService, AgentsReader: agentStore,
Sessions: executionStore,
Sessions: sessionService,
SessionsReader: sessionStore,
SessionCreation: executionStore,
SessionEvents: executionStore,
Turns: executionStore,
SessionEvents: sessionStore,
Turns: sessionStore,
Items: sessionStore,
Subagents: sessionStore,
Artifacts: sessionService,
ArtifactsReader: sessionStore,
SessionAdmin: executionStore,
SessionAdmin: sessionStore,
Environments: sessionService, EnvironmentsReader: sessionStore, ExecutorConnections: executorConnections{sessions: sessionStore, registry: registry},
Admin: executionStore, AdminAudit: auditStore, WriteAudit: auditStore, Metrics: metrics,
Admin: sessionStore, AdminAudit: auditStore, WriteAudit: auditStore, Metrics: metrics,
RuntimeObservations: observationService, RuntimeHistory: historyService,
}
if worker != nil {
deps.Execution = &api.Execution{
ExecutorURL: executorURL,
SessionAdmission: worker,
InputAdmission: worker,
SessionArchive: worker,
SessionArchive: deploymentExecution,
Workspaces: worker,
NativeInstaller: nativeInstaller,
}
Expand Down
6 changes: 3 additions & 3 deletions services/core/cmd/server/runtime_history.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,11 +12,11 @@ import (
"time"

"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/coremetrics"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimehistory"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimehistory/postgresreader"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs/otlpexporter"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store"
"golang.org/x/net/http/httpguts"
)

Expand Down Expand Up @@ -52,7 +52,7 @@ type runtimeHistoryExporter interface {
Close(context.Context) error
}

func runtimeHistory(ctx context.Context, coreStore *store.Store, executionEnabled bool) (runtimeHistorySetup, error) {
func runtimeHistory(ctx context.Context, units *pgunit.Pool, executionEnabled bool) (runtimeHistorySetup, error) {
config, err := loadRuntimeHistoryConfig()
if err != nil {
return runtimeHistorySetup{}, err
Expand All @@ -70,7 +70,7 @@ func runtimeHistory(ctx context.Context, coreStore *store.Store, executionEnable
Metrics: []runtimehistory.Metric{runtimehistory.MetricCPU, runtimehistory.MetricMemory, runtimehistory.MetricTokens},
}
timeout := time.Duration(config.TimeoutSeconds) * time.Second
backend, err := postgresreader.New(coreStore, postgresreader.Config{Capabilities: capabilities, QueryTimeout: timeout})
backend, err := postgresreader.New(units, postgresreader.Config{Capabilities: capabilities, QueryTimeout: timeout})
if err != nil {
return runtimeHistorySetup{}, err
}
Expand Down
6 changes: 3 additions & 3 deletions services/core/cmd/server/runtime_history_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,8 @@ package main

import (
"context"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeobs"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store"
"os"
"path/filepath"
"strings"
Expand All @@ -19,7 +19,7 @@ func (panicHistoryExporter) Close(context.Context) error
func TestRuntimeHistoryUsesCoreDatabaseByDefault(t *testing.T) {
t.Setenv("OAC_HISTORY_SETTINGS_FILE", "")
for _, enabled := range []bool{true, false} {
setup, err := runtimeHistory(t.Context(), store.New(nil), enabled)
setup, err := runtimeHistory(t.Context(), pgunit.NewPool(nil), enabled)
if err != nil {
t.Fatal(err)
}
Expand All @@ -45,7 +45,7 @@ func TestRuntimeHistoryOptionalExportAndSamplingConfiguration(t *testing.T) {
t.Fatal(err)
}
t.Setenv("OAC_HISTORY_SETTINGS_FILE", file)
setup, err := runtimeHistory(t.Context(), store.New(nil), true)
setup, err := runtimeHistory(t.Context(), pgunit.NewPool(nil), true)
if err != nil {
t.Fatal(err)
}
Expand Down
8 changes: 3 additions & 5 deletions services/core/internal/api/admin_resources.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,19 +5,17 @@ import (
"net/http"

"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store"
"github.com/go-chi/chi/v5"
)

// adminTenantContextKey identifies an explicit management target, not a caller.
// Only Core-key-authenticated project resource handlers receive it.
type adminTenantContextKey struct{}

// Admin reads the administrator's cross-Project views: the asset summary and
// the Sessions whose Runtime is observed.
// Admin reads the administrator's cross-Project Session views.
type Admin interface {
ReadAdminSummary(context.Context, string, store.AdminSummaryFilter, func(sessions.Session, *string) error) (store.AdminAssetCounts, error)
ListAdminRuntimeTargets(context.Context, []string, string, int, bool) (store.AdminRuntimeTargetPage, error)
ReadAdminSummary(context.Context, string, sessions.AdminSummaryFilter, func(sessions.Session, *string) error) (sessions.AdminAssetCounts, error)
ListAdminRuntimeTargets(context.Context, []string, string, int, bool) (sessions.AdminRuntimeTargetPage, error)
}

func (h *Handler) adminResourceScope(next http.Handler) http.Handler {
Expand Down
9 changes: 4 additions & 5 deletions services/core/internal/api/admin_resources_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,6 @@ import (
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/identity"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/projects"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store"
)

const managementProjectID = "22222222-2222-4222-8222-222222222222"
Expand Down Expand Up @@ -117,21 +116,21 @@ func TestAdminResourcesHaveExplicitTargetWithoutCallerImpersonation(t *testing.T

type summaryFixture struct {
tenant string
filter store.AdminSummaryFilter
filter sessions.AdminSummaryFilter
}

func (s *summaryFixture) ReadAdminSummary(_ context.Context, tenant string, filter store.AdminSummaryFilter, visit func(sessions.Session, *string) error) (store.AdminAssetCounts, error) {
func (s *summaryFixture) ReadAdminSummary(_ context.Context, tenant string, filter sessions.AdminSummaryFilter, visit func(sessions.Session, *string) error) (sessions.AdminAssetCounts, error) {
s.tenant, s.filter = tenant, filter
for i, usage := range []json.RawMessage{nil, json.RawMessage(`{"input_tokens":3,"output_tokens":5,"total_tokens":8,"input_tokens_details":{"cached_tokens":2},"output_tokens_details":{"reasoning_tokens":1}}`)} {
session := sessions.Session{ID: "session", TenantID: tenant, Configuration: json.RawMessage(`{"agent":{"id":"agent","model":"model","tools":[]},"environment":{"type":"none"}}`), CreatedAt: time.Unix(100+int64(i), 0), Usage: usage}
if i == 0 {
session.LastTurn = &sessions.Turn{Status: sessions.TurnInProgress, CreatedAt: time.Unix(110, 0)}
}
if err := visit(session, nil); err != nil {
return store.AdminAssetCounts{}, err
return sessions.AdminAssetCounts{}, err
}
}
return store.AdminAssetCounts{Agents: 4, Skills: 2}, nil
return sessions.AdminAssetCounts{Agents: 4, Skills: 2}, nil
}
func TestAdminSummaryUsesPublicStateAndNullUsageCoverage(t *testing.T) {
key := callerBinding()
Expand Down
2 changes: 1 addition & 1 deletion services/core/internal/api/admin_runtime.go
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,7 @@ func (h *Handler) adminRuntimeObservations(w http.ResponseWriter, r *http.Reques
}
page, err := h.Admin.ListAdminRuntimeTargets(ctx, tenants, options.after, options.limit, options.ascending)
if err != nil {
writeStoreError(w, r, err)
writeSessionsError(w, r, err)
return
}
sessions := make([]runtimeobs.SessionIdentity, len(page.Data))
Expand Down
Loading
Loading