diff --git a/services/core/tests/integration/environment_expiry_worker_test.go b/services/core/tests/integration/environment_expiry_worker_test.go index 04838a170..a10b81a36 100644 --- a/services/core/tests/integration/environment_expiry_worker_test.go +++ b/services/core/tests/integration/environment_expiry_worker_test.go @@ -9,6 +9,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/sessions" "github.com/google/uuid" @@ -110,7 +111,9 @@ func TestWorkerEnvironmentExpiryWithoutDevicesAndAfterRestart(t *testing.T) { if err != nil || got.State != sessions.EnvironmentInputPending || !got.Deadline.Equal(future.Deadline) { t.Fatal("future input changed", got, err) } + awaitRelease := pgtest.ObserveExecutionLeaseRelease(t, s.pool) stop() + awaitRelease() makeEnvironmentExpiryDue(t, s.pool, &future) _, stop = startEnvironmentExpiryWorker(t, s, d) waitEnvironmentExpiry(t, s, futureTenant, future) diff --git a/services/core/tests/integration/environment_worker_test.go b/services/core/tests/integration/environment_worker_test.go index b6cfbad55..af6326bc6 100644 --- a/services/core/tests/integration/environment_worker_test.go +++ b/services/core/tests/integration/environment_worker_test.go @@ -6,6 +6,7 @@ import ( "time" "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/sessions" ) @@ -155,7 +156,9 @@ func TestWorkerEnvironmentRetriesPendingWithoutExtendingDeadline(t *testing.T) { t.Fatal("active preparation was duplicated", frame.Type) case <-time.After(time.Second): } + awaitRelease := pgtest.ObserveExecutionLeaseRelease(t, h.s.pool) stop() + awaitRelease() nextWorkerFrame(t, frames, proto.TypeExecutionRelease) stored, err := sessionAdapter(h.s).GetEnvironmentInputReservation(t.Context(), h.tenant, pending.SessionID, pending.ID) if err != nil || stored.State != sessions.EnvironmentInputPending || !stored.Deadline.Equal(pending.Deadline) || len(stored.Receipts) != 0 {