diff --git a/contracts/agents-api/runtime-observability.md b/contracts/agents-api/runtime-observability.md index f1a961d79..07c94b355 100644 --- a/contracts/agents-api/runtime-observability.md +++ b/contracts/agents-api/runtime-observability.md @@ -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. + +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. diff --git a/contracts/agents-api/zh/runtime-observability.md b/contracts/agents-api/zh/runtime-observability.md index a6d894a45..8743b7648 100644 --- a/contracts/agents-api/zh/runtime-observability.md +++ b/contracts/agents-api/zh/runtime-observability.md @@ -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`。 @@ -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 失败日志标识退出阶段。失去执行租约仍会终止执行;不会自动重放查询或外部执行。 diff --git a/deploy/kubernetes/BETA_CHANGELOG.md b/deploy/kubernetes/BETA_CHANGELOG.md index 085753133..800a37bf6 100644 --- a/deploy/kubernetes/BETA_CHANGELOG.md +++ b/deploy/kubernetes/BETA_CHANGELOG.md @@ -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. diff --git a/services/core/cmd/server/main.go b/services/core/cmd/server/main.go index 044305694..9f4d77ba8 100644 --- a/services/core/cmd/server/main.go +++ b/services/core/cmd/server/main.go @@ -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) } } @@ -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) } diff --git a/services/core/cmd/server/runtime_history_ownership.go b/services/core/cmd/server/runtime_history_ownership.go new file mode 100644 index 000000000..378b59916 --- /dev/null +++ b/services/core/cmd/server/runtime_history_ownership.go @@ -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 +} diff --git a/services/core/cmd/server/runtime_history_ownership_test.go b/services/core/cmd/server/runtime_history_ownership_test.go new file mode 100644 index 000000000..8eaa2c35f --- /dev/null +++ b/services/core/cmd/server/runtime_history_ownership_test.go @@ -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", ¬Owned, "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") + } + }) + } +} diff --git a/services/core/internal/execution/worker.go b/services/core/internal/execution/worker.go index 4e9a901dd..6783ddd1d 100644 --- a/services/core/internal/execution/worker.go +++ b/services/core/internal/execution/worker.go @@ -17,6 +17,8 @@ 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 @@ -24,6 +26,8 @@ type Worker struct { dispatcher *Dispatcher admission *store.Store lease Ownership + ownershipCheckOnce sync.Once + ownershipChecks chan struct{} directoryReads chan directoryReadRequest fileWrites chan fileWriteRequest scheduleWake chan struct{} @@ -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 @@ -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() @@ -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() { @@ -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 { @@ -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 } } @@ -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) diff --git a/services/core/internal/execution/worker_failure.go b/services/core/internal/execution/worker_failure.go new file mode 100644 index 000000000..fe8fae066 --- /dev/null +++ b/services/core/internal/execution/worker_failure.go @@ -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)) +} diff --git a/services/core/internal/execution/worker_failure_test.go b/services/core/internal/execution/worker_failure_test.go new file mode 100644 index 000000000..85401c6ff --- /dev/null +++ b/services/core/internal/execution/worker_failure_test.go @@ -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()) + } +} diff --git a/services/core/internal/execution/worker_metrics.go b/services/core/internal/execution/worker_metrics.go index c549452ca..cb2337d3a 100644 --- a/services/core/internal/execution/worker_metrics.go +++ b/services/core/internal/execution/worker_metrics.go @@ -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" diff --git a/services/core/internal/execution/worker_metrics_test.go b/services/core/internal/execution/worker_metrics_test.go index c8dcb374f..eb57397d1 100644 --- a/services/core/internal/execution/worker_metrics_test.go +++ b/services/core/internal/execution/worker_metrics_test.go @@ -5,6 +5,7 @@ import ( "errors" "fmt" "sync" + "sync/atomic" "testing" "time" @@ -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") + } +} diff --git a/services/core/internal/persistence/postgres/pgunit/lease.go b/services/core/internal/persistence/postgres/pgunit/lease.go index 392472d02..8d7355256 100644 --- a/services/core/internal/persistence/postgres/pgunit/lease.go +++ b/services/core/internal/persistence/postgres/pgunit/lease.go @@ -55,7 +55,7 @@ func AcquireLease(ctx context.Context, pool *pgxpool.Pool) (*Lease, error) { // connection and commits only when apply returns nil. apply receives the // context carrying the execution deadline. func (l *Lease) Transaction(ctx context.Context, apply func(context.Context, pgx.Tx) error) error { - return l.withConn(ctx, func(ctx context.Context, conn *pgxpool.Conn) error { + return l.withConn(ctx, "transaction", func(ctx context.Context, conn *pgxpool.Conn) error { return run(ctx, conn, readWrite, apply) }) } @@ -63,7 +63,7 @@ func (l *Lease) Transaction(ctx context.Context, apply func(context.Context, pgx // CheckOwnership pings the leased connection, confirming that this service // still owns the database before external work. func (l *Lease) CheckOwnership(ctx context.Context) error { - return l.withConn(ctx, func(ctx context.Context, conn *pgxpool.Conn) error { return conn.Ping(ctx) }) + return l.withConn(ctx, "ownership_check", func(ctx context.Context, conn *pgxpool.Conn) error { return conn.Ping(ctx) }) } // CancelOperations cancels coordinator-owned contexts between leased @@ -76,7 +76,7 @@ func (l *Lease) CancelOperations(ctx context.Context, cancel context.CancelFunc) if cancel == nil { return errors.New("execution lease cancellation requires a cancel function") } - return l.withConn(ctx, func(ctx context.Context, conn *pgxpool.Conn) error { + return l.withConn(ctx, "cancellation_fence", func(ctx context.Context, conn *pgxpool.Conn) error { if err := conn.Ping(ctx); err != nil { return err } @@ -114,13 +114,20 @@ func (l *Lease) Close(ctx context.Context) error { } // withConn applies the execution deadline to the gate wait and the operation. -func (l *Lease) withConn(ctx context.Context, apply func(context.Context, *pgxpool.Conn) error) error { - ctx, cancel := context.WithTimeout(ctx, ExecutionTimeout) +func (l *Lease) withConn(caller context.Context, operation string, apply func(context.Context, *pgxpool.Conn) error) (err error) { + started := time.Now() + ctx, cancel := context.WithTimeout(caller, ExecutionTimeout) defer cancel() - if err := l.lock(ctx); err != nil { + if err = l.lock(ctx); err != nil { + observeLeaseFailure(caller, ctx, operation, "gate", started, err, nil) return err } defer l.unlock() + defer func() { + if err != nil { + observeLeaseFailure(caller, ctx, operation, "connection", started, err, l.conn) + } + }() if l.conn == nil { return ErrLeaseClosed } diff --git a/services/core/internal/persistence/postgres/pgunit/lease_diagnostics.go b/services/core/internal/persistence/postgres/pgunit/lease_diagnostics.go new file mode 100644 index 000000000..73059280e --- /dev/null +++ b/services/core/internal/persistence/postgres/pgunit/lease_diagnostics.go @@ -0,0 +1,50 @@ +package pgunit + +import ( + "context" + "errors" + "fmt" + "net" + "time" + + obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" + "github.com/jackc/pgx/v5/pgconn" + "github.com/jackc/pgx/v5/pgxpool" +) + +// Never log error text: PostgreSQL errors can contain SQL and confidential data. +func leaseErrorClass(err error) string { + switch { + case err == nil: + return "none" + case errors.Is(err, context.Canceled): + return "cancelled" + case errors.Is(err, context.DeadlineExceeded): + return "deadline_exceeded" + case errors.Is(err, pgconn.ErrConnClosed): + return "connection_closed" + case errors.Is(err, ErrLeaseClosed): + return "lease_closed" + } + var pgError *pgconn.PgError + if errors.As(err, &pgError) { + return "postgres_error" + } + var networkError net.Error + if errors.As(err, &networkError) { + return "network_error" + } + return "operation_error" +} + +func observeLeaseFailure(caller, operationContext context.Context, operation, phase string, started time.Time, err error, conn *pgxpool.Conn) { + fields := []any{"operation", operation, "phase", phase, "duration_ms", float64(time.Since(started)) / float64(time.Millisecond), "error_class", leaseErrorClass(err), "error_type", fmt.Sprintf("%T", err), "caller_context", leaseErrorClass(caller.Err()), "operation_context", leaseErrorClass(operationContext.Err())} + if phase == "connection" { + fields = append(fields, "connection_closed", conn == nil || conn.Conn().PgConn().IsClosed()) + } + var pgError *pgconn.PgError + if errors.As(err, &pgError) { + fields = append(fields, "sqlstate", pgError.Code) + } + obslog.Ctx(caller).Warn("execution lease operation failed", fields...) +} diff --git a/services/core/internal/persistence/postgres/pgunit/lease_diagnostics_test.go b/services/core/internal/persistence/postgres/pgunit/lease_diagnostics_test.go new file mode 100644 index 000000000..477ecc278 --- /dev/null +++ b/services/core/internal/persistence/postgres/pgunit/lease_diagnostics_test.go @@ -0,0 +1,33 @@ +package pgunit + +import ( + "bytes" + "context" + "errors" + "log/slog" + "strings" + "testing" + "time" + + "github.com/jackc/pgx/v5/pgconn" +) + +func TestLeaseDiagnosticsRetainClassWithoutConfidentialErrorText(t *testing.T) { + var output bytes.Buffer + previous := slog.Default() + slog.SetDefault(slog.New(slog.NewJSONHandler(&output, nil))) + t.Cleanup(func() { slog.SetDefault(previous) }) + err := &pgconn.PgError{Code: "57014", Message: "secret-canary", Detail: "secret-canary", Where: "secret-canary", InternalQuery: "secret-canary"} + observeLeaseFailure(t.Context(), t.Context(), "transaction", "connection", time.Now(), err, nil) + if strings.Contains(output.String(), "secret-canary") || !strings.Contains(output.String(), `"sqlstate":"57014"`) || !strings.Contains(output.String(), `"connection_closed":true`) { + t.Fatal("unsafe or incomplete lease failure diagnostic", output.String()) + } + for _, tt := range []struct { + err error + want string + }{{nil, "none"}, {context.Canceled, "cancelled"}, {context.DeadlineExceeded, "deadline_exceeded"}, {pgconn.ErrConnClosed, "connection_closed"}, {ErrLeaseClosed, "lease_closed"}, {err, "postgres_error"}, {errors.New("secret-canary"), "operation_error"}} { + if got := leaseErrorClass(tt.err); got != tt.want { + t.Fatalf("class=%s want=%s", got, tt.want) + } + } +} diff --git a/services/core/internal/persistence/postgres/pgunit/lease_sampling_cancellation_test.go b/services/core/internal/persistence/postgres/pgunit/lease_sampling_cancellation_test.go new file mode 100644 index 000000000..48afad5cb --- /dev/null +++ b/services/core/internal/persistence/postgres/pgunit/lease_sampling_cancellation_test.go @@ -0,0 +1,89 @@ +package pgunit + +import ( + "context" + "errors" + "net" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" + "github.com/jackc/pgx/v5/pgxpool" +) + +type samplingDelayConn struct { + net.Conn + armed *atomic.Bool + entered chan struct{} + release <-chan struct{} + deadline chan struct{} + signal *sync.Once +} + +func (c *samplingDelayConn) Read(p []byte) (int, error) { + if c.armed.CompareAndSwap(true, false) { + close(c.entered) + <-c.release + } + return c.Conn.Read(p) +} + +func (c *samplingDelayConn) SetDeadline(t time.Time) error { + if !t.IsZero() && t.Before(time.Now().Add(50*time.Millisecond)) { + c.signal.Do(func() { close(c.deadline) }) + } + return c.Conn.SetDeadline(t) +} + +// A short observational context can destroy the real writer's connection even +// though Ping writes no application data. Samplers must never pass such contexts +// to the lease; only the execution owner performs authoritative checks. +func TestCancelledOwnershipPingClosesExecutionLease(t *testing.T) { + var armed atomic.Bool + entered, release := make(chan struct{}), make(chan struct{}) + deadline := make(chan struct{}) + var signal, unblock sync.Once + releaseRead := func() { unblock.Do(func() { close(release) }) } + pool := pgtest.OpenIsolated(t, func(cfg *pgxpool.Config) { + dial := cfg.ConnConfig.DialFunc + cfg.ConnConfig.DialFunc = func(ctx context.Context, network, address string) (net.Conn, error) { + conn, err := dial(ctx, network, address) + if err != nil { + return nil, err + } + return &samplingDelayConn{Conn: conn, armed: &armed, entered: entered, release: release, deadline: deadline, signal: &signal}, nil + } + }) + lease := acquireLease(t, pool) + t.Cleanup(releaseRead) + armed.Store(true) + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + done := make(chan error, 1) + go func() { done <- lease.CheckOwnership(ctx) }() + select { + case <-entered: + case <-time.After(time.Second): + t.Fatal("ownership Ping did not read") + } + cancel() + select { + case <-deadline: + case <-time.After(time.Second): + t.Fatal("cancellation did not reach the driver deadline") + } + releaseRead() + select { + case err := <-done: + if !errors.Is(err, context.Canceled) || !lease.conn.Conn().PgConn().IsClosed() { + t.Fatal("cancelled Ping did not invalidate lease", err) + } + case <-time.After(2 * time.Second): + t.Fatal("cancelled Ping did not settle") + } + if lease.CheckOwnership(t.Context()) == nil { + t.Fatal("cancelled observational Ping retained writer ownership") + } +} diff --git a/services/core/internal/runtimeobs/sampler.go b/services/core/internal/runtimeobs/sampler.go index 2fa81163f..1360c7364 100644 --- a/services/core/internal/runtimeobs/sampler.go +++ b/services/core/internal/runtimeobs/sampler.go @@ -33,6 +33,8 @@ type HistoryObserver interface { ObserveSessionsForHistory(context.Context, []SessionIdentity, OwnershipChecker, PageOptions) ([]Observation, []error) } +// OwnershipChecker observes permission to sample, never execution authority. +// Its short-lived contexts must not operate on the execution lease connection. type OwnershipChecker interface { CheckOwnership(context.Context) error }