From 776a9414f2d1eaf2ee01eec94a40c29791650b10 Mon Sep 17 00:00:00 2001 From: saladday <1203511142@qq.com> Date: Thu, 1 Oct 2026 02:12:52 +0800 Subject: [PATCH] Observe server-side lease release before test ownership handoff --- services/core/IMPLEMENTATION.md | 2 +- .../postgres/pgtest/lease.go} | 21 +++++++++++++------ .../store/environment_claim_worker_test.go | 5 +++-- .../environment_connection_worker_test.go | 3 ++- .../store/environment_initial_input_test.go | 3 +++ .../store/environment_input_activity_test.go | 3 +++ .../store/environment_input_migration_test.go | 3 +++ .../internal/store/environment_inputs_test.go | 3 +++ .../store/environment_write_audit_test.go | 3 +++ .../store/runtime_allocations_test.go | 3 +++ .../runtime_enrollment_connection_test.go | 3 ++- .../store/sandbox_deployment_setup_test.go | 3 +++ .../store/subagent_identities_test.go | 3 +++ 13 files changed, 47 insertions(+), 11 deletions(-) rename services/core/internal/{store/execution_lease_handoff_test.go => persistence/postgres/pgtest/lease.go} (73%) diff --git a/services/core/IMPLEMENTATION.md b/services/core/IMPLEMENTATION.md index d6adda08b..1e3064844 100644 --- a/services/core/IMPLEMENTATION.md +++ b/services/core/IMPLEMENTATION.md @@ -175,7 +175,7 @@ Enabling the daemon gateway with `OAC_PUBLIC_URL` also starts the execution Work Sandbox reset snapshots bind the deployment relation explicitly to its single row before joining resources, so that even on a fresh database without statistics an inflated join estimate cannot trigger JIT compilation inside the lease deadline. -At startup the Worker fails previously claimed work, keeps queued input and never replays uncertain execution. Shutdown cancels active dispatch and attempts terminal persistence before releasing the lease; a lost owner cannot commit. Closing the lease invalidates its writer and waits for pgx cleanup within the caller's deadline; a later close can resume that wait. Tests that transfer ownership immediately must observe the previous owner's advisory lock disappear before starting the next, with a bounded wait that fails on query errors. The lease fences database writes, not already queued daemon commands or native effects. +At startup the Worker fails previously claimed work, keeps queued input and never replays uncertain execution. Shutdown cancels active dispatch and attempts terminal persistence before releasing the lease; a lost owner cannot commit. Closing the lease invalidates its writer and waits for pgx cleanup within the caller's deadline; a later close can resume that wait. Tests that transfer ownership immediately use `pgtest.ObserveExecutionLeaseRelease` to observe the previous owner's advisory lock disappear before starting the next, with a bounded wait that fails on query errors. The lease fences database writes, not already queued daemon commands or native effects. ## Hosted sandboxes diff --git a/services/core/internal/store/execution_lease_handoff_test.go b/services/core/internal/persistence/postgres/pgtest/lease.go similarity index 73% rename from services/core/internal/store/execution_lease_handoff_test.go rename to services/core/internal/persistence/postgres/pgtest/lease.go index ffbcf500f..551c16f22 100644 --- a/services/core/internal/store/execution_lease_handoff_test.go +++ b/services/core/internal/persistence/postgres/pgtest/lease.go @@ -1,4 +1,4 @@ -package store_test +package pgtest import ( "context" @@ -8,9 +8,9 @@ import ( "github.com/jackc/pgx/v5/pgxpool" ) -// observeExecutionLeaseRelease captures the current owner before shutdown. Local +// ObserveExecutionLeaseRelease captures the current owner before shutdown. Local // pgx cleanup does not acknowledge the server's release of its advisory lock. -func observeExecutionLeaseRelease(t *testing.T, pool *pgxpool.Pool) func() { +func ObserveExecutionLeaseRelease(t *testing.T, pool *pgxpool.Pool) func() { t.Helper() // Match the single-bigint key in queries/scheduling.sql, scoped to this DB. const lock = `locktype='advisory' AND granted AND objsubid=1 @@ -26,12 +26,21 @@ func observeExecutionLeaseRelease(t *testing.T, pool *pgxpool.Pool) func() { t.Helper() ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second) defer cancel() - awaitDaemonRemoteCondition(t, ctx, 5*time.Second, "previous execution lease release", func() bool { + ticker := time.NewTicker(10 * time.Millisecond) + defer ticker.Stop() + for { var held bool if err := pool.QueryRow(ctx, "SELECT EXISTS (SELECT 1 FROM pg_locks WHERE "+lock+" AND pid=$1)", owner).Scan(&held); err != nil { t.Fatal("observe execution lease release", err) } - return !held - }) + if !held { + return + } + select { + case <-ctx.Done(): + t.Fatal("previous execution lease was not released", ctx.Err()) + case <-ticker.C: + } + } } } diff --git a/services/core/internal/store/environment_claim_worker_test.go b/services/core/internal/store/environment_claim_worker_test.go index b55853550..dbde88971 100644 --- a/services/core/internal/store/environment_claim_worker_test.go +++ b/services/core/internal/store/environment_claim_worker_test.go @@ -7,6 +7,7 @@ import ( "testing" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store" @@ -37,7 +38,7 @@ func TestWorkerReconcilesEnvironmentPromotionBeforeStart(t *testing.T) { t.Fatal("promotion did not retain the active claim", turn, err) } // Simulate owner loss after commit, without sending any daemon Start. - awaitRelease := observeExecutionLeaseRelease(t, db.pool) + awaitRelease := pgtest.ObserveExecutionLeaseRelease(t, db.pool) if err := owner.Lease.Close(t.Context()); err != nil { t.Fatal(err) } @@ -45,7 +46,7 @@ func TestWorkerReconcilesEnvironmentPromotionBeforeStart(t *testing.T) { restarted := startWorker(t, t.Context(), db, &execution.Dispatcher{Store: s, Registry: runtimegateway.NewRegistry()}) stopped, cancel := context.WithCancel(t.Context()) cancel() - awaitRelease = observeExecutionLeaseRelease(t, db.pool) + awaitRelease = pgtest.ObserveExecutionLeaseRelease(t, db.pool) if err := restarted.Run(stopped); !errors.Is(err, context.Canceled) { t.Fatal(err) } diff --git a/services/core/internal/store/environment_connection_worker_test.go b/services/core/internal/store/environment_connection_worker_test.go index c7c9ed8a0..df47a830e 100644 --- a/services/core/internal/store/environment_connection_worker_test.go +++ b/services/core/internal/store/environment_connection_worker_test.go @@ -7,6 +7,7 @@ import ( "time" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store" "github.com/google/uuid" @@ -32,7 +33,7 @@ func TestEnvironmentConnectionWorkerReconcilesAndReleasesLease(t *testing.T) { if err := writer.ObserveEnvironmentConnection(t.Context(), tenant, environment.ID, generation, 1, true); err != nil { t.Fatal(err) } - awaitRelease := observeExecutionLeaseRelease(t, db.pool) + awaitRelease := pgtest.ObserveExecutionLeaseRelease(t, db.pool) if err := owner.Lease.Close(t.Context()); err != nil { t.Fatal(err) } diff --git a/services/core/internal/store/environment_initial_input_test.go b/services/core/internal/store/environment_initial_input_test.go index 025dbc046..830a59385 100644 --- a/services/core/internal/store/environment_initial_input_test.go +++ b/services/core/internal/store/environment_initial_input_test.go @@ -8,6 +8,7 @@ import ( "testing" "time" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" "github.com/google/uuid" "github.com/jackc/pgx/v5" @@ -228,9 +229,11 @@ func TestEnvironmentInitialInputExpiryHasNoTurnAndCannotReplay(t *testing.T) { if err != nil || len(events) != expectedEvents || events[len(events)-1].Event.Type != "agent.session.failed" || events[len(events)-1].Turn != nil || events[len(events)-1].EnvironmentInputActivity.Status != "failed" { t.Fatal("missing pre-Turn failure snapshot", events, err) } + awaitRelease := pgtest.ObserveExecutionLeaseRelease(t, writer.pool) if err := writer.lease.Close(t.Context()); err != nil { t.Fatal(err) } + awaitRelease() pool.Close() reopened, reopenedPool := testStore(t) writer = executionWriter(t, reopened) diff --git a/services/core/internal/store/environment_input_activity_test.go b/services/core/internal/store/environment_input_activity_test.go index 8fb106e3e..e0a414608 100644 --- a/services/core/internal/store/environment_input_activity_test.go +++ b/services/core/internal/store/environment_input_activity_test.go @@ -6,6 +6,7 @@ import ( "reflect" "testing" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" "github.com/google/uuid" "github.com/jackc/pgx/v5" @@ -216,9 +217,11 @@ func TestEnvironmentInputActivityRecoversWaitingActionAndHidesDeletion(t *testin t.Fatal(err) } requireEnvironmentInputActivity(t, s, tenant, session.ID, "idle", "") + awaitRelease := pgtest.ObserveExecutionLeaseRelease(t, old.pool) if err := old.lease.Close(t.Context()); err != nil { t.Fatal(err) } + awaitRelease() next := executionWriter(t, s) if err := next.ReconcileEnvironmentConnections(t.Context()); err != nil { t.Fatal(err) diff --git a/services/core/internal/store/environment_input_migration_test.go b/services/core/internal/store/environment_input_migration_test.go index 68c587430..6b8815c80 100644 --- a/services/core/internal/store/environment_input_migration_test.go +++ b/services/core/internal/store/environment_input_migration_test.go @@ -7,6 +7,7 @@ import ( "strings" "testing" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" "github.com/google/uuid" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/stdlib" @@ -99,9 +100,11 @@ func TestEnvironmentInputPromotionUsesCurrentExecutionWriter(t *testing.T) { t.Fatal("pooled Store promoted input without execution ownership") } closed := executionWriter(t, s) + awaitRelease := pgtest.ObserveExecutionLeaseRelease(t, closed.pool) if err := closed.lease.Close(t.Context()); err != nil { t.Fatal(err) } + awaitRelease() if _, err := closed.PromoteEnvironmentInput(t.Context(), tenant, session.ID, pending.ID); err == nil { t.Fatal("closed execution writer promoted pending input") } diff --git a/services/core/internal/store/environment_inputs_test.go b/services/core/internal/store/environment_inputs_test.go index e49368ca2..9b31a4f6e 100644 --- a/services/core/internal/store/environment_inputs_test.go +++ b/services/core/internal/store/environment_inputs_test.go @@ -9,6 +9,7 @@ import ( "testing" "time" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" "github.com/google/uuid" "github.com/jackc/pgx/v5/pgxpool" @@ -166,9 +167,11 @@ func TestEnvironmentInputReservationPromotionAndDirectRetries(t *testing.T) { if err != nil || len(retry) != 2 || !retry[0].Replayed || retry[0].Sequence != promoted.Receipts[0].Sequence { t.Fatal("direct retry after promotion", retry, err) } + awaitRelease := pgtest.ObserveExecutionLeaseRelease(t, writer.pool) if err := writer.lease.Close(ctx); err != nil { t.Fatal(err) } + awaitRelease() pool.Close() restarted, pool := testStore(t) after, err := executionWriter(t, restarted).PromoteEnvironmentInput(ctx, tenant, session.ID, first.ID) diff --git a/services/core/internal/store/environment_write_audit_test.go b/services/core/internal/store/environment_write_audit_test.go index eaee546bd..8c38d9578 100644 --- a/services/core/internal/store/environment_write_audit_test.go +++ b/services/core/internal/store/environment_write_audit_test.go @@ -8,6 +8,7 @@ import ( "sync" "testing" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/writeaudit" ) @@ -22,9 +23,11 @@ func TestEnvironmentUploadWriteAuditSurvivesRequestAndLeaseContext(t *testing.T) cancel() sessionAuditCount(t, f.s, f.tenant, 0) // Neither the settlement caller nor the newly acquired execution lease owns the request context. + awaitRelease := pgtest.ObserveExecutionLeaseRelease(t, f.writer.pool) if err := f.writer.lease.Close(context.Background()); err != nil { t.Fatal(err) } + awaitRelease() next := executionWriter(t, f.s) var group sync.WaitGroup for range 4 { diff --git a/services/core/internal/store/runtime_allocations_test.go b/services/core/internal/store/runtime_allocations_test.go index f25ac5480..700f9e2e7 100644 --- a/services/core/internal/store/runtime_allocations_test.go +++ b/services/core/internal/store/runtime_allocations_test.go @@ -7,6 +7,7 @@ import ( "testing" "time" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" "github.com/google/uuid" ) @@ -31,9 +32,11 @@ func TestRuntimeAllocationAtomicOwnershipAndRecovery(t *testing.T) { if _, err := w.ReserveRuntimeAllocation(t.Context(), tenant, environment.ID, uuid.NewString(), runtimedevice.HashCredential(secret)); !errors.Is(err, ErrIdempotencyConflict) { t.Fatalf("provider target changed: %v", err) } + awaitRelease := pgtest.ObserveExecutionLeaseRelease(t, w.pool) if err := w.lease.Close(t.Context()); err != nil { t.Fatal(err) } + awaitRelease() if _, err := w.ObserveRuntimeRunning(t.Context(), owner); err == nil { t.Fatal("lost writer changed allocation") } diff --git a/services/core/internal/store/runtime_enrollment_connection_test.go b/services/core/internal/store/runtime_enrollment_connection_test.go index 0f8806716..3441b2562 100644 --- a/services/core/internal/store/runtime_enrollment_connection_test.go +++ b/services/core/internal/store/runtime_enrollment_connection_test.go @@ -13,6 +13,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtime" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeenrollment" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store" @@ -145,7 +146,7 @@ func TestEnrolledDaemonConnectionRevocationAndRestart(t *testing.T) { second := connect(rotated.Token) await("connected") assertConnection(environment.ID, rotated.Token, "connected", 200) - awaitRelease := observeExecutionLeaseRelease(t, db.pool) + awaitRelease := pgtest.ObserveExecutionLeaseRelease(t, db.pool) stop() stop = nil awaitRelease() diff --git a/services/core/internal/store/sandbox_deployment_setup_test.go b/services/core/internal/store/sandbox_deployment_setup_test.go index 15620df86..73f20b911 100644 --- a/services/core/internal/store/sandbox_deployment_setup_test.go +++ b/services/core/internal/store/sandbox_deployment_setup_test.go @@ -8,6 +8,7 @@ import ( "sync" "testing" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox/e2b" "github.com/google/uuid" @@ -70,9 +71,11 @@ func TestSandboxDeploymentSetupPersistsWithoutExecution(t *testing.T) { if err := pool.QueryRow(t.Context(), "SELECT (SELECT count(*) FROM sessions)+(SELECT count(*) FROM runtime_nodes)+(SELECT count(*) FROM runtime_allocations)+(SELECT count(*) FROM runtime_placements)").Scan(&sideEffects); err != nil || sideEffects != 0 { t.Fatal("setup or rejected admission created execution state", sideEffects, err) } + awaitRelease := pgtest.ObserveExecutionLeaseRelease(t, w.pool) if err := w.lease.Close(context.Background()); err != nil { t.Fatal(err) } + awaitRelease() restarted := executionWriter(t, s) if err := restarted.ClaimWebSandboxDeployment(t.Context(), id); err != nil { t.Fatal(err) diff --git a/services/core/internal/store/subagent_identities_test.go b/services/core/internal/store/subagent_identities_test.go index a6d5ce057..26277bbfb 100644 --- a/services/core/internal/store/subagent_identities_test.go +++ b/services/core/internal/store/subagent_identities_test.go @@ -8,6 +8,7 @@ import ( "testing" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" "github.com/google/uuid" @@ -110,9 +111,11 @@ func TestSubagentIdentityIsAtomicScopedAndImmutable(t *testing.T) { if _, err = w.CompleteExecution(ctx, tenant, session.ID, input.TurnID, sessions.TurnCompleted, json.RawMessage(`{}`), "root", input.Sequence); err != nil { t.Fatal(err) } + awaitRelease := pgtest.ObserveExecutionLeaseRelease(t, w.pool) if err = w.lease.Close(ctx); err != nil { t.Fatal(err) } + awaitRelease() reopened, _ := testStore(t) nextOwner := executionWriter(t, reopened) again, err := reopened.GetSubagentIdentity(ctx, tenant, session.ID, "child-a")