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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions contracts/agents-api/runtime-observability.md
Original file line number Diff line number Diff line change
Expand Up @@ -128,3 +128,9 @@ With an OTLP endpoint configured, Core exports every record, both `on_read` and
CPU and memory points are exported only when the sample has `started_at`; a missing measurement produces no point. Attributes are `agents.tenant.id`, `agents.session.id`, `agents.environment.id`, `agents.runtime.allocation.id`, `agents.runtime.mode`, `agents.runtime.provider.type`, `agents.runtime.status`, `agents.runtime.reason`, `agents.runtime.collection.source` and nanosecond `agents.runtime.resolved_at_unix_nano`, `agents.runtime.observed_at_unix_nano` and `agents.runtime.compute.started_at_unix_nano`. The nanosecond times keep records joinable when a backend stores event time at lower precision. Provider keys, receipts, native identifiers, raw errors, paths and credentials are never attributes.

Web reads history only through Core; neither a Collector nor another metrics store is needed for its charts.

## Sampling ownership

Periodic history sampling reads the execution Worker's observed ownership state. Its short source deadlines and sweep cancellation never run database operations on the connection holding the execution lease. The Worker retains its authoritative database ownership checks before execution, publishes their observations in check order, and invalidates the observed state when it stops or loses ownership. Waiting for an ownership check and performing it share one bounded deadline; a caller cancellation while waiting never starts a database operation. An unknown, failed or stopped ownership observation prevents sampling; it never grants execution authority.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Reconcile the existing sampling contract with observed ownership

This paragraph says periodic sampling only reads the Worker's observed ownership state, but the same document still promises that collection runs under the database lease, checks that lease every 100 ms (lines 88–90), and performs node-history copying after a lease check (line 106); the Chinese mirror retains the same conflicting claims at lines 90–92 and 108. Since the implementation now consults runtimeHistoryOwnership rather than the lease, update those existing paragraphs instead of leaving contradictory operational semantics.

AGENTS.md reference: AGENTS.md:L58-L58

Useful? React with 👍 / 👎.


Failed leased operations log their operation category, gate or connection phase, error class and type, caller and operation cancellation state, connection-closed state when observed under the gate, duration and PostgreSQL SQLSTATE when present. These diagnostics omit error text, SQL, credentials and model data. Worker failure logs identify the exiting stage. A lost execution lease remains fatal; no query or external execution is automatically replayed.
8 changes: 7 additions & 1 deletion 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: "9d5cfbb9046979ec8e1e9d730ce35618e968698769571242c0aea0596abd9858"
---

这是面向贡献者的契约,规定 Core 如何观测 Runtime 并保留其历史。路由和响应字段见 [Runtime telemetry API](runtime-observability-api.md)。代码位于 `services/core/internal/runtimeobs`(解析、源、采样器和导出)、`internal/runtimehistory`(历史查询和 PostgreSQL 存储)以及 `internal/runtimeobs/otlpexporter`。
Expand Down Expand Up @@ -130,3 +130,9 @@ PostgreSQL 存储仅保留周期性的 `openai_hosted` 记录,因此 API 读
仅当采样包含 `started_at` 时才导出 CPU 和内存数据点;缺少测量值不会产生数据点。属性包括 `agents.tenant.id`、`agents.session.id`、`agents.environment.id`、`agents.runtime.allocation.id`、`agents.runtime.mode`、`agents.runtime.provider.type`、`agents.runtime.status`、`agents.runtime.reason`、`agents.runtime.collection.source`,以及以纳秒为单位的 `agents.runtime.resolved_at_unix_nano`、`agents.runtime.observed_at_unix_nano` 和 `agents.runtime.compute.started_at_unix_nano`。这些纳秒时间使记录在后端以较低精度存储事件时间时仍可关联。Provider key、回执、原生标识符、原始错误、路径和凭据绝不会作为属性。

Web 仅通过 Core 读取历史记录;其图表不需要 Collector 或其他指标存储。

## 采样所有权 {#sampling-ownership}

周期历史采样读取执行 Worker 已观测到的所有权状态。采样的短来源超时和整轮取消不会在持有执行租约的连接上执行数据库操作。Worker 在执行前保留权威的数据库所有权检查,按检查顺序发布观测结果,并在停止或失去所有权时将观测状态标为不可用。等待所有权检查和执行检查共用一个有界期限;调用方在等待时取消,不会启动数据库操作。未知、失败或已停止的所有权观测会阻止采样,永远不授予执行权限。

失败的租约操作记录操作类别、等待连接闸门或使用连接的阶段、错误类别和类型、调用方与操作的取消状态、在持有闸门时观测到的连接关闭状态、耗时,以及存在时的 PostgreSQL SQLSTATE。诊断不记录错误原文、SQL、凭据或模型数据。Worker 失败日志标识退出阶段。失去执行租约仍会终止执行;不会自动重放查询或外部执行。
12 changes: 12 additions & 0 deletions deploy/kubernetes/BETA_CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -64,3 +64,15 @@ The administrative deployment update records authenticated administrator provena
A combined image retains all packaged Harnesses, while each managed Runtime discovers and registers only the Session-selected Harness. Unknown or unavailable selections fail without substituting another implementation. Self-hosted installations retain their installed Harness set. Codex version probes still validate the executable and expose secret-safe process spawn/wait timing.

Publish matching Core, provider helpers and a newly built Runtime template together. Node deployments require matching protocol-version-6 nodes. Existing allocations retain their bootstrap and Runtime; qualification must use a fresh allocation. [Runtime bootstrap](../../docs/runtime-bootstrap.md) owns the startup contract.

## 2026-10-07 — Isolate history sampling from execution ownership

| Item | Value |
| --- | --- |
| Fork beta baseline | `c68a68443d1856fcde3f64130b559b331de12184` |
| Feature branch | `codex/core-lease-failure-diagnostics` |
| Schema migrations / DDL / SQL changes | None |
| Runtime bootstrap / provider helpers / node wire / native pins | Unchanged |
| Required release assets | Matching Core/Web release; retain the selected Runtime template |

History sampling observes Worker ownership without using its short-lived contexts on the execution lease connection. Execution retains authoritative lease checks and fails closed on ownership loss. Worker shutdown immediately invalidates the sampling observation. Lease failure logs identify cancellation, timeout and connection state without SQL or error text, and Worker failure logs identify the exiting stage. [Runtime observability](../../contracts/agents-api/runtime-observability.md#sampling-ownership) owns these rules.
7 changes: 4 additions & 3 deletions services/core/cmd/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -81,7 +81,7 @@ func main() {
return
}
if err := run(); err != nil {
log.Bg().Error("oac-core startup failed", "error", err)
log.Bg().Error("oac-core failed", "error", err)
os.Exit(1)
}
}
Expand Down Expand Up @@ -350,11 +350,12 @@ func run() error {
if worker == nil {
return errors.New("Runtime history periodic sampling requires the execution worker")
}
sampler, err := runtimeobs.NewSampler(observationResolver, observationService, worker, runtimeobs.SamplerOptions{
historyOwner := runtimeHistoryOwnership{snapshot: worker.MetricsSnapshot}
sampler, err := runtimeobs.NewSampler(observationResolver, observationService, historyOwner, runtimeobs.SamplerOptions{
Interval: history.SampleInterval,
Report: func(result runtimeobs.SweepResult) {
sampleCtx, cancel := context.WithTimeout(ctx, 2*time.Second)
sampleErr := worker.CheckOwnership(sampleCtx)
sampleErr := historyOwner.CheckOwnership(sampleCtx)
if sampleErr == nil {
_, sampleErr = deploymentStore.SampleHostHistory(sampleCtx)
}
Expand Down
26 changes: 26 additions & 0 deletions services/core/cmd/server/runtime_history_ownership.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
package main

import (
"context"
"errors"

"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution"
)

// runtimeHistoryOwnership implements the sampler's read-only ownership check.
// The execution worker performs authoritative database checks. Sampling observes
// their result without putting its short-lived contexts on the leased connection.
type runtimeHistoryOwnership struct {
snapshot func() execution.WorkerMetrics
}

func (o runtimeHistoryOwnership) CheckOwnership(ctx context.Context) error {
if err := ctx.Err(); err != nil {
return err
}
snapshot := o.snapshot()
if snapshot.ExecutionOwner == nil || !*snapshot.ExecutionOwner || snapshot.Scheduler.Status == "stopped" || snapshot.Scheduler.Status == "failing" {
return errors.New("execution ownership is not observed")
}
return nil
}
41 changes: 41 additions & 0 deletions services/core/cmd/server/runtime_history_ownership_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
package main

import (
"context"
"testing"

"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution"
)

func TestHistoryOwnershipObservesWorkerWithoutDatabaseChecks(t *testing.T) {
owned, notOwned := true, false
for _, tt := range []struct {
name string
owner *bool
status string
available bool
}{
{"unobserved", nil, "unknown", false}, {"owned", &owned, "ok", true},
{"acquired before first poll", &owned, "unknown", true}, {"released", &notOwned, "stopped", false},
{"draining", &owned, "stopped", false}, {"failed", &owned, "failing", false},
} {
t.Run(tt.name, func(t *testing.T) {
calls := 0
owner := runtimeHistoryOwnership{snapshot: func() execution.WorkerMetrics {
calls++
return execution.WorkerMetrics{ExecutionOwner: tt.owner, Scheduler: execution.WorkerJobMetrics{Status: tt.status}}
}}
if got := owner.CheckOwnership(t.Context()) == nil; got != tt.available {
t.Fatalf("ownership availability = %v", got)
}
if calls != 1 {
t.Fatal("snapshot was not read once")
}
ctx, cancel := context.WithCancel(t.Context())
cancel()
if owner.CheckOwnership(ctx) == nil || calls != 1 {
t.Fatal("cancelled sampling touched the owner")
}
})
}
}
31 changes: 30 additions & 1 deletion services/core/internal/execution/worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,13 +17,17 @@ import (

const DefaultExecutionConcurrency = 4

const ownershipCheckTimeout = 5 * time.Second

// Worker owns queued work; the database lease excludes a second execution service.
type Worker struct {
concurrency int
metrics workerMetricsState
dispatcher *Dispatcher
admission *store.Store
lease Ownership
ownershipCheckOnce sync.Once
ownershipChecks chan struct{}
directoryReads chan directoryReadRequest
fileWrites chan fileWriteRequest
scheduleWake chan struct{}
Expand Down Expand Up @@ -120,6 +124,20 @@ func StartWorker(ctx context.Context, dispatcher *Dispatcher, owner Owner) (_ *W

// CheckOwnership checks the same database lease used for execution writes.
func (w *Worker) CheckOwnership(ctx context.Context) error {
ctx, cancel := context.WithTimeout(ctx, ownershipCheckTimeout)
defer cancel()
w.ownershipCheckOnce.Do(func() { w.ownershipChecks = make(chan struct{}, 1) })
select {
case w.ownershipChecks <- struct{}{}:
case <-ctx.Done():
return ctx.Err()
}
defer func() { <-w.ownershipChecks }()
if err := ctx.Err(); err != nil {
return err
}
// Keep the authoritative result and its observation in the same order.
// A delayed success must not overwrite a later failed ownership check.
err := w.lease.CheckOwnership(ctx)
w.observeOwnership(err)
return err
Expand Down Expand Up @@ -172,9 +190,13 @@ func (w *Worker) CreateSessionStream(ctx context.Context, tenant string, input s
func (w *Worker) Run(ctx context.Context) (runErr error) {
defer w.stopOnce.Do(func() { close(w.stopped) })
ctx, cancel := context.WithCancel(ctx)
exitStage := "context"
var running sync.WaitGroup
defer func() {
w.observeWorkerStop(runErr, ctx.Err())
if runErr != nil && !errors.Is(runErr, ctx.Err()) {
observeWorkerFailure(ctx, exitStage, runErr)
}
cancel()
if w.runtimes != nil {
w.runtimes.stop()
Expand Down Expand Up @@ -229,8 +251,10 @@ func (w *Worker) Run(ctx context.Context) (runErr error) {
case <-ctx.Done():
return ctx.Err()
case err := <-preparationDone:
exitStage = "environment_initialization"
return err
case err := <-lifecycleDone:
exitStage = "runtime_lifecycle"
return err
case request := <-w.fileWrites:
if request.ctx.Err() != nil || active[request.environment.SessionID] || len(active) == w.executionConcurrency() {
Expand Down Expand Up @@ -290,6 +314,7 @@ func (w *Worker) Run(ctx context.Context) (runErr error) {
delete(active, result.id)
w.observeSlots(len(active))
if result.err != nil {
exitStage = "execution_completion"
return result.err
}
if !rescanOnCompletion {
Expand All @@ -303,20 +328,23 @@ func (w *Worker) Run(ctx context.Context) (runErr error) {
case <-ticker.C:
maintenance = true
}
check, stop := context.WithTimeout(ctx, 5*time.Second)
check, stop := context.WithTimeout(ctx, ownershipCheckTimeout)
err := w.CheckOwnership(check)
stop()
if err != nil {
w.observeSchedulerPoll(0, err)
exitStage = "ownership_check"
return err
}
if maintenance {
if _, err := w.dispatcher.Store.ExpireEnvironmentInputs(ctx); err != nil {
w.observeSchedulerPoll(0, err)
exitStage = "expire_environment_inputs"
return err
}
if err := w.observeEnrolledRuntimes(ctx); err != nil {
w.observeSchedulerPoll(0, err)
exitStage = "observe_enrolled_runtimes"
return err
}
}
Expand All @@ -333,6 +361,7 @@ func (w *Worker) Run(ctx context.Context) (runErr error) {
w.observeSlots(len(active))
if err != nil {
w.observeSchedulerPoll(0, err)
exitStage = "select_work"
return err
}
w.observeSchedulerPoll(len(work), nil)
Expand Down
12 changes: 12 additions & 0 deletions services/core/internal/execution/worker_failure.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
package execution

import (
"context"
"fmt"

obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log"
)

func observeWorkerFailure(ctx context.Context, stage string, err error) {
obslog.Ctx(ctx).Error("execution worker stopped", "stage", stage, "error_type", fmt.Sprintf("%T", err))
}
25 changes: 25 additions & 0 deletions services/core/internal/execution/worker_failure_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
package execution

import (
"bytes"
"encoding/json"
"errors"
"log/slog"
"strings"
"testing"
)

func TestWorkerFailureLogPreservesCompletionStageWithoutErrorText(t *testing.T) {
var output bytes.Buffer
previous := slog.Default()
slog.SetDefault(slog.New(slog.NewJSONHandler(&output, nil)))
t.Cleanup(func() { slog.SetDefault(previous) })
observeWorkerFailure(t.Context(), "execution_completion", errors.New("confidential-query-canary"))
var event map[string]any
if err := json.Unmarshal(output.Bytes(), &event); err != nil {
t.Fatal(err)
}
if event["stage"] != "execution_completion" || event["error_type"] == nil || strings.Contains(output.String(), "confidential-query-canary") {
t.Fatal("completion failure lost provenance or leaked error text", output.String())
}
}
2 changes: 2 additions & 0 deletions services/core/internal/execution/worker_metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,8 @@ func (w *Worker) observeSchedulerPoll(processed int, err error) {
func (w *Worker) observeWorkerStop(runErr, contextErr error) {
w.metrics.mu.Lock()
defer w.metrics.mu.Unlock()
w.metrics.closed = true
w.metrics.value.ExecutionOwner = nil
w.metrics.value.Scheduler.Status = "stopped"
if runErr != nil && (contextErr == nil || !errors.Is(runErr, contextErr)) {
w.metrics.value.Scheduler.Status = "failing"
Expand Down
74 changes: 74 additions & 0 deletions services/core/internal/execution/worker_metrics_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"errors"
"fmt"
"sync"
"sync/atomic"
"testing"
"time"

Expand Down Expand Up @@ -107,3 +108,76 @@ func TestWorkerMetricsConcurrentSnapshots(t *testing.T) {
}
running.Wait()
}

func TestWorkerStopInvalidatesOwnershipBeforeLeaseDrain(t *testing.T) {
worker := &Worker{}
worker.observeOwnership(nil)
worker.observeWorkerStop(context.Canceled, context.Canceled)
worker.observeOwnership(nil)
snapshot := worker.MetricsSnapshot()
if snapshot.ExecutionOwner != nil || snapshot.Scheduler.Status != "stopped" {
t.Fatal("late check revived ownership while draining")
}
worker.observeWorkerClosed(nil)
if owner := worker.MetricsSnapshot().ExecutionOwner; owner == nil || *owner {
t.Fatal("confirmed release was not recorded")
}
}

type orderedOwnershipLease struct {
calls atomic.Int32
entered chan struct{}
release <-chan struct{}
}

func (l *orderedOwnershipLease) CheckOwnership(context.Context) error {
if l.calls.Add(1) == 1 {
close(l.entered)
<-l.release
return nil
}
return pgunit.ErrLeaseClosed
}
func (*orderedOwnershipLease) CancelOperations(context.Context, context.CancelFunc) error { return nil }
func (*orderedOwnershipLease) Close(context.Context) error { return nil }

func TestWorkerOwnershipChecksPublishInOrderAndRespectWaitingDeadline(t *testing.T) {
release := make(chan struct{})
var unblock sync.Once
lease := &orderedOwnershipLease{entered: make(chan struct{}), release: release}
worker := &Worker{lease: lease}
first := make(chan error, 1)
go func() { first <- worker.CheckOwnership(t.Context()) }()
t.Cleanup(func() {
unblock.Do(func() { close(release) })
select {
case <-first:
case <-time.After(time.Second):
t.Error("first ownership check did not stop")
}
})
<-lease.entered
waiting, cancel := context.WithTimeout(t.Context(), 25*time.Millisecond)
defer cancel()
if err := worker.CheckOwnership(waiting); !errors.Is(err, context.DeadlineExceeded) {
t.Fatal("waiting check ran ahead of the authoritative result", err)
}
if lease.calls.Load() != 1 {
t.Fatal("waiting check reached the lease out of order")
}
unblock.Do(func() { close(release) })
if err := <-first; err != nil {
t.Fatal(err)
}
// Leave a result for cleanup after the successful first caller has settled.
first <- nil
if owner := worker.MetricsSnapshot().ExecutionOwner; owner == nil || !*owner {
t.Fatal("successful check was not observed")
}
if err := worker.CheckOwnership(t.Context()); !errors.Is(err, pgunit.ErrLeaseClosed) {
t.Fatal(err)
}
if worker.MetricsSnapshot().ExecutionOwner != nil {
t.Fatal("older success revived a failed ownership observation")
}
}
Loading