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
2 changes: 1 addition & 1 deletion services/core/IMPLEMENTATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
package store_test
package pgtest

import (
"context"
Expand All @@ -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
Expand All @@ -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:
}
}
}
}
5 changes: 3 additions & 2 deletions services/core/internal/store/environment_claim_worker_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -37,15 +38,15 @@ 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)
}
awaitRelease()
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)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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")
}
Expand Down
3 changes: 3 additions & 0 deletions services/core/internal/store/environment_inputs_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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)
Expand Down
3 changes: 3 additions & 0 deletions services/core/internal/store/environment_write_audit_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)
Expand All @@ -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 {
Expand Down
3 changes: 3 additions & 0 deletions services/core/internal/store/runtime_allocations_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)
Expand All @@ -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")
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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()
Expand Down
3 changes: 3 additions & 0 deletions services/core/internal/store/sandbox_deployment_setup_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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)
Expand Down
3 changes: 3 additions & 0 deletions services/core/internal/store/subagent_identities_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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")
Expand Down