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 CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -110,7 +110,7 @@ Run `make check` before completing code changes. The `check` target in the [Make
| `OAC_TEST_DATABASE_URL` | A dedicated test database. The full gate fails when it is missing. |
| `OAC_TEST_OFFICIAL_SDK_PYTHON` | The pinned official SDK interpreter |

The role needs `CREATE DATABASE`: managed-provider tests create and drop isolated `oac_*_tests` databases because provider identity is deployment-wide. Tests must not bypass the production provider-switch guard.
The role needs `CREATE DATABASE`: tests of database-wide state, such as the execution lease and the provider identity, create and drop isolated `oac_*_tests` databases. Tests must not bypass the production provider-switch guard.

### Contract and schema rules

Expand Down
8 changes: 7 additions & 1 deletion services/core/IMPLEMENTATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,12 @@

These are the code-level rules of `services/core` that no contract states. Contracts own behavior: the [coverage ledger](../../contracts/agents-api/README.md) and its linked contracts own the public and administrator APIs, the [machine connection API](../../contracts/agents-api/machine-api.md) the `/api/v1` routes, the [Core–Runtime protocol](../../docs/runtime-protocol.md) the daemon wire, and the [Sandbox Provider guide](../../docs/sandbox-provider.md#managed-lifecycle) the managed compute lifecycle. When code changes one of these rules, change the rule here in the same branch.

## Layering

`internal/persistence/postgres/pgunit` owns Core's PostgreSQL transaction and execution-lease mechanics: pooled read-write and snapshot transactions, the lease's dedicated connection and its gate, the ownership check, the cancellation fence, close, and the execution deadline. Persistence code runs every transaction through it, and nothing outside `persistence` and `store` imports it. `internal/persistence/postgres/pgtest` is test support: it opens the dedicated test database under the `oac_*_tests` guard, applies the migrations, and creates isolated databases for database-wide state such as the execution lease. Only test files import it.

`store` is transitional. `store.New` builds a pooled Store, and `store.NewExecution` takes the lease and builds the execution writer on it. An execution-only operation on a pooled Store fails with `store.ErrExecutionAuthority`. New adapters do not copy that check: their execution repositories require a `*pgunit.Lease` at construction, their public repositories expose no execution operation, and the check goes away with `store`.

## Request handling

Every Agents API JSON route reads its body through `readJSONObject` before decoding, validation or lookup. The gate requires a JSON Content-Type, applies the route's body limit and rejects invalid UTF-8, malformed JSON (including unpaired surrogate escapes), repeated keys and non-object roots with the official messages; an empty body or `null` becomes `{}`. DELETE, multipart, Core extension and internal routes keep their own readers. Member names match exactly: decode request objects with `decodeInputObject`, or check `inexactMember` before another decoder, so `encoding/json` never matches a case variant.
Expand Down Expand Up @@ -143,7 +149,7 @@ Item merging never mutates the incoming observation or the previous snapshot: pu

## Worker ownership

Enabling the daemon gateway with `OAC_PUBLIC_URL` also starts the execution Worker. One Worker owns an execution database through a dedicated PostgreSQL advisory-lock connection, and its store view uses that connection for every Session transaction: binding, claim and reconciliation, journal, Items and usage, function callbacks and receipts, and terminal state. These short transactions and the lease pings serialize, with a five-second deadline that includes gate and Session-lock waits. Never hold a transaction across daemon or model work, reconnect the writer or fall back to the pool after losing the lease. Public admission and device maintenance use pooled connections, and pooled reads grant no write authority.
Enabling the daemon gateway with `OAC_PUBLIC_URL` also starts the execution Worker. One Worker owns an execution database through a `pgunit.Lease`: the PostgreSQL advisory lock held by one dedicated connection. A second Worker on the same database fails to start. The Worker's execution writer runs every Session transaction on that connection: binding, claim and reconciliation, journal, Items and usage, function callbacks and receipts, and terminal state. The lease gate serializes these short transactions and the ownership pings, and `pgunit.ExecutionTimeout` (five seconds) bounds each one, including its gate and Session-lock waits. Cancelling an in-flight pgx operation can close the connection that owns the lock, so coordinator-owned work is cancelled only through the lease's cancellation fence, between leased operations. Never hold a transaction across daemon or model work, reconnect the writer or fall back to the pool after losing the lease. Public admission and device maintenance use pooled connections, and pooled reads grant no write authority. Transactions state their isolation: read committed for writes, read-only repeatable read for snapshots.

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.

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -62,8 +62,7 @@ func TestArchiveWaitingCleanupReceiptBarrier(t *testing.T) {
}{{"Kill_no_delivery", false, false}, {"KillCompute_no_delivery", true, false}, {"Kill_live_delivery", false, true}, {"KillCompute_live_delivery", true, true}} {
t.Run(scenario.name, func(t *testing.T) {
checkpoint := scenario.checkpoint
s, lease := resetManagerStore(t)
writer := lease.Store()
s, writer := resetManagerStore(t)
installation := uuid.NewString()
if err := writer.ClaimWebSandboxDeployment(t.Context(), installation); err != nil {
t.Fatal(err)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ type finishObservationFixture struct {
func newFinishObservationFixture(t *testing.T, maxConnections int32) finishObservationFixture {
t.Helper()
var cfg *pgxpool.Config
s, lease := resetManagerStoreConfig(t, func(c *pgxpool.Config) {
s, writer := resetManagerStoreConfig(t, func(c *pgxpool.Config) {
if maxConnections > 0 {
c.MaxConns = maxConnections
}
Expand Down Expand Up @@ -56,7 +56,7 @@ func newFinishObservationFixture(t *testing.T, maxConnections int32) finishObser
if err != nil {
t.Fatal(err)
}
return finishObservationFixture{s, lease.Store(), pool, tenant, session, Dispatcher{Store: lease.Store()}}
return finishObservationFixture{s, writer, pool, tenant, session, Dispatcher{Store: writer}}
}
func (f finishObservationFixture) start(t *testing.T) store.InputReceipt {
t.Helper()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,82 +3,14 @@ package execution
import (
"context"
"errors"
"net"
"os"
"strings"
"sync"
"sync/atomic"
"testing"
"time"

"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store"
"github.com/google/uuid"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)

// Each case has its own real advisory-lock namespace. The delayed driver read
// lets the cancellation fence hit its own deadline without shortening production
// timeouts or depending on an unrelated SQL failure to exercise this branch.
func retirementFailureLease(t *testing.T, armed *atomic.Bool, reading chan struct{}, release <-chan struct{}) (*store.ExecutionLease, *pgxpool.Pool) {
t.Helper()
dsn := os.Getenv("OAC_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("dedicated PostgreSQL required")
}
cfg, err := pgxpool.ParseConfig(dsn)
if err != nil {
t.Fatal(err)
}
if !strings.HasPrefix(cfg.ConnConfig.Database, "oac_") || !strings.HasSuffix(cfg.ConnConfig.Database, "_tests") {
t.Fatal("dedicated test database required")
}
admin, err := pgxpool.NewWithConfig(t.Context(), cfg.Copy())
if err != nil {
t.Fatal(err)
}
t.Cleanup(admin.Close)
name := "oac_retirement_" + uuid.NewString()[:8] + "_tests"
quoted := pgx.Identifier{name}.Sanitize()
if _, err := admin.Exec(t.Context(), "CREATE DATABASE "+quoted); err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
if _, err := admin.Exec(ctx, "DROP DATABASE "+quoted+" WITH (FORCE)"); err != nil {
t.Error(err)
}
})
cfg.ConnConfig.Database = name
dial := cfg.ConnConfig.DialFunc
cfg.ConnConfig.DialFunc = func(ctx context.Context, network, address string) (net.Conn, error) {
c, err := dial(ctx, network, address)
if err != nil {
return nil, err
}
return &delayedLeaseRead{Conn: c, armed: armed, reading: reading, release: release}, nil
}
pool, err := pgxpool.NewWithConfig(t.Context(), cfg)
if err != nil {
t.Fatal(err)
}
t.Cleanup(pool.Close)
lease, err := store.New(pool).AcquireExecutionLease(t.Context())
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
if err := lease.Close(ctx); err != nil {
t.Error(err)
}
})
return lease, pool
}

func TestFailedInventoryRetirementClosesAdmissionAndRetainsGate(t *testing.T) {
for _, mode := range []string{"gate_timeout", "lease_loss"} {
t.Run(mode, func(t *testing.T) {
Expand All @@ -87,9 +19,9 @@ func TestFailedInventoryRetirementClosesAdmissionAndRetainsGate(t *testing.T) {
var readOnce sync.Once
unblockRead := func() { readOnce.Do(func() { close(releaseRead) }) }
defer unblockRead()
lease, pool := retirementFailureLease(t, &armed, reading, releaseRead)
writer, pool := delayedReadWriter(t, &armed, reading, releaseRead)
m := testRuntimeManager(t)
m.store = lease.Store()
m.store = writer
m.loadDeployment = func(context.Context) (*RuntimeProvider, error) { return nil, nil }
m.mutationGate = make(chan struct{}, 1)
// This fixture models an already loaded node deployment; its provider is
Expand Down Expand Up @@ -117,7 +49,7 @@ func TestFailedInventoryRetirementClosesAdmissionAndRetainsGate(t *testing.T) {
defer cancel()
queryDone = make(chan error, 1)
armed.Store(true)
go func() { queryDone <- lease.Ping(queryCtx) }()
go func() { queryDone <- writer.CheckExecutionOwnership(queryCtx) }()
select {
case <-reading:
case <-time.After(2 * time.Second):
Expand Down
93 changes: 38 additions & 55 deletions services/core/internal/execution/sandbox_deployment_drain_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,18 +4,17 @@ import (
"bytes"
"context"
"net"
"os"
"strings"
"sync"
"sync/atomic"
"testing"
"time"

"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/sandbox/node"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store"
"github.com/google/uuid"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)

Expand Down Expand Up @@ -43,6 +42,35 @@ func (c *delayedLeaseRead) Read(p []byte) (int, error) {
return c.Conn.Read(p)
}

// delayedReadWriter owns an isolated database, so each case has its own
// advisory-lock namespace. The delayed driver read lets a cancellation fence hit
// its own deadline without shortening production timeouts.
func delayedReadWriter(t *testing.T, armed *atomic.Bool, reading chan struct{}, release <-chan struct{}) (*store.Store, *pgxpool.Pool) {
t.Helper()
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) {
c, err := dial(ctx, network, address)
if err != nil {
return nil, err
}
return &delayedLeaseRead{Conn: c, armed: armed, reading: reading, release: release}, nil
}
})
writer, err := store.NewExecution(t.Context(), store.New(pool))
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
if err := writer.CloseExecution(ctx); err != nil {
t.Error(err)
}
})
return writer, pool
}

func TestSandboxDeploymentDrainPreservesLeaseInFlightRead(t *testing.T) {
for _, mode := range []string{"deployment", "inventory", "manual"} {
t.Run(mode, func(t *testing.T) {
Expand All @@ -52,62 +80,17 @@ func TestSandboxDeploymentDrainPreservesLeaseInFlightRead(t *testing.T) {
}

func testLifecycleCancellationPreservesLease(t *testing.T, mode string) {
dsn := os.Getenv("OAC_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("dedicated PostgreSQL required")
}
cfg, err := pgxpool.ParseConfig(dsn)
if err != nil {
t.Fatal(err)
}
if !strings.HasPrefix(cfg.ConnConfig.Database, "oac_") || !strings.HasSuffix(cfg.ConnConfig.Database, "_tests") {
t.Fatal("dedicated test database required")
}
admin, err := pgxpool.NewWithConfig(t.Context(), cfg.Copy())
if err != nil {
t.Fatal(err)
}
defer admin.Close()
database := pgx.Identifier{"oac_drain_" + uuid.NewString()[:8] + "_tests"}.Sanitize()
if _, err := admin.Exec(t.Context(), "CREATE DATABASE "+database); err != nil {
t.Fatal(err)
}
defer func() {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
if _, err := admin.Exec(ctx, "DROP DATABASE "+database+" WITH (FORCE)"); err != nil {
t.Error(err)
}
}()
cfg.ConnConfig.Database = strings.Trim(database, `"`)
var armed atomic.Bool
reading, release := make(chan struct{}), make(chan struct{})
var releaseOnce sync.Once
unblock := func() { releaseOnce.Do(func() { close(release) }) }
defer unblock()
dial := cfg.ConnConfig.DialFunc
cfg.ConnConfig.DialFunc = func(ctx context.Context, network, address string) (net.Conn, error) {
c, err := dial(ctx, network, address)
if err != nil {
return nil, err
}
return &delayedLeaseRead{Conn: c, armed: &armed, reading: reading, release: release}, nil
}
pool, err := pgxpool.NewWithConfig(t.Context(), cfg)
if err != nil {
t.Fatal(err)
}
defer pool.Close()
lease, err := store.New(pool).AcquireExecutionLease(t.Context())
if err != nil {
t.Fatal(err)
}
defer lease.Close(context.Background())
writer, _ := delayedReadWriter(t, &armed, reading, release)
hub := node.NewHub(node.HubOptions{})
defer hub.Close()
id := uuid.NewString()
configuration := &RuntimeProvider{InstallationID: id, ProviderKind: "docker", Mode: "nodes", Generation: 1, CoreURL: "https://core.example/api/v1", BackendFingerprint: strings.Repeat("a", 64), Provider: hub.Proxy(uuid.NewString(), "docker", 1)}
m, err := newRuntimeManager(lease.Store(), runtimegateway.NewRegistry(), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return configuration, nil }))
m, err := newRuntimeManager(writer, runtimegateway.NewRegistry(), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return configuration, nil }))
if err != nil {
t.Fatal(err)
}
Expand Down Expand Up @@ -141,11 +124,11 @@ func testLifecycleCancellationPreservesLease(t *testing.T, mode string) {
return
}
defer finish()
queryDone <- lease.Store().CheckExecutionOwnership(ctx)
queryDone <- writer.CheckExecutionOwnership(ctx)
<-ctx.Done()
// Provider settlement remains outside the cancellation fence. A fresh owner
// read must proceed even before this tracked lifecycle operation returns.
leaseFree <- lease.Store().CheckExecutionOwnership(t.Context())
leaseFree <- writer.CheckExecutionOwnership(t.Context())
}()
select {
case <-reading:
Expand Down Expand Up @@ -198,7 +181,7 @@ func testLifecycleCancellationPreservesLease(t *testing.T, mode string) {
if delayed != nil {
delayed.unblock()
}
if err := lease.Store().CheckExecutionOwnership(t.Context()); err != nil {
if err := writer.CheckExecutionOwnership(t.Context()); err != nil {
t.Fatalf("deployment drain destroyed the owner connection (in-flight query: %v): %v", queryErr, err)
}
if queryErr != nil {
Expand Down Expand Up @@ -240,12 +223,12 @@ func (c *delayedCancellationContext) unblock() {
}

func TestSandboxDeploymentDrainFailureCannotReactivate(t *testing.T) {
_, lease := resetManagerStore(t)
_, writer := resetManagerStore(t)
hub := node.NewHub(node.HubOptions{})
defer hub.Close()
id := uuid.NewString()
configuration := &RuntimeProvider{InstallationID: id, ProviderKind: "docker", Mode: "nodes", Generation: 1, CoreURL: "https://core.example/api/v1", BackendFingerprint: strings.Repeat("a", 64), Provider: hub.Proxy(uuid.NewString(), "docker", 1)}
m, err := newRuntimeManager(lease.Store(), runtimegateway.NewRegistry(), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return configuration, nil }))
m, err := newRuntimeManager(writer, runtimegateway.NewRegistry(), NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return configuration, nil }))
if err != nil {
t.Fatal(err)
}
Expand All @@ -257,7 +240,7 @@ func TestSandboxDeploymentDrainFailureCannotReactivate(t *testing.T) {
if err != nil {
t.Fatal(err)
}
if err := lease.Close(t.Context()); err != nil {
if err := writer.CloseExecution(t.Context()); err != nil {
t.Fatal(err)
}
first := m.pauseDeployment(t.Context())
Expand Down
3 changes: 1 addition & 2 deletions services/core/internal/execution/sandbox_generations_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,8 +15,7 @@ import (
)

func TestE2BReplacementVerifiesTwiceAndNeverPublishesFailedCommit(t *testing.T) {
s, lease := resetManagerStore(t)
writer := lease.Store()
s, writer := resetManagerStore(t)
id := uuid.NewString()
if err := writer.ClaimWebSandboxDeployment(t.Context(), id); err != nil {
t.Fatal(err)
Expand Down
Loading
Loading