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
8 changes: 5 additions & 3 deletions services/core/IMPLEMENTATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,11 +6,13 @@ These are the code-level rules of `services/core` that no contract states. Contr

`api.NewHandler` takes one `api.Dependencies` value, built only in `cmd/server`. Each application area is one field typed as an interface declared in `api` beside its handlers, listing exactly the methods they call. Every field is required and `NewHandler` rejects a missing one, except the optional groups whose comments say what nil means: `Execution` is nil without an execution Worker, `Sandboxes` is nil without a managed sandbox installation and requires `Execution`, and `Execution.NativeInstaller` is nil for a build without a source revision. Handlers never discover a capability by type assertion or fall back to another implementation. API tests use one strict fake per area, `fake<Area>`, which fails the test on any call the test did not set.

`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.
`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. Outside `persistence` and `store`, only `cmd/server`, which acquires the lease, and test fixtures import 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.

Shared vocabulary has one owner each, and domains use it rather than copy it. `internal/environmentconfig` owns Environment setup, Skills, Plugins and initial files with their validation and public metadata; `Setup.Validate` checks requested configuration, where a Skill may be an unresolved reference, and `Setup.ValidateInstalled` checks frozen, installable configuration. `internal/skills` owns `ParseVersion`, the canonical positive decimal Skill version. `internal/metadata` owns the metadata rules: `Validate` for the pair, key and value limits and U+0000, `ValidateStorable` for U+0000 alone, and `Encode` with its 64 KiB bound. `internal/jsonobject` owns `Normalize`, the stable encoding of stored JSON objects that snapshots and retry identities compare. These packages import no persistence.

`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`.
`store` is transitional. `store.New` builds a pooled Store, and `store.NewExecution(s, lease)` builds the execution writer on a lease it borrows. 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`.

`cmd/server` owns the execution lease. It acquires one `pgunit.Lease`, builds every lease-bound adapter on it, and passes the lease and those adapters together as one `execution.Owner` to `execution.StartWorker`. If anything fails before that call, `cmd/server` closes the lease. From that call the Worker owns cleanup: a failed start closes the lease before it returns, and a started Worker closes it after `Run` has cancelled and drained its work. Each close runs under its own bounded deadline, independent of the cancelled request or run. Lease-bound adapters and the store writer borrow the lease and never close it, and the Worker uses the lease only through `Owner.Lease`, never through an adapter. Store integration tests start the Worker the same way through `startWorker`.

## Request handling

Expand Down Expand Up @@ -153,7 +155,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 `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.
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 Core on the same database cannot acquire the lease and starts no Worker. 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
9 changes: 7 additions & 2 deletions services/core/cmd/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ import (
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/databaseurl"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/nativeinstaller"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit"
"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/runtimegateway"
Expand Down Expand Up @@ -226,8 +227,12 @@ func run() error {
if registry != nil {
dispatcher := &execution.Dispatcher{Store: executionStore, Registry: registry,
ManagedRuntimes: managed, MaxConcurrentExecutions: concurrency}

worker, err = execution.StartWorker(ctx, dispatcher)
lease, err := pgunit.AcquireLease(ctx, pool)
if err != nil {
return err
}
// From this call on the Worker closes the lease, even when it fails to start.
worker, err = execution.StartWorker(ctx, dispatcher, execution.Owner{Lease: lease, Store: store.NewExecution(executionStore, lease)})
if err != nil {
return err
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,8 @@ 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, writer := resetManagerStore(t)
s, leased := resetManagerStore(t)
writer := leased.Store
installation := uuid.NewString()
if err := writer.ClaimWebSandboxDeployment(t.Context(), installation); err != nil {
t.Fatal(err)
Expand Down Expand Up @@ -167,7 +168,7 @@ func TestArchiveWaitingCleanupReceiptBarrier(t *testing.T) {
t.Fatal("Kill bypassed durable cleanup ownership", allocation, err)
}
}}
lifecycle := &runtimeLifecycle{store: writer, registry: registry, config: RuntimeProvider{InstallationID: installation, Provider: provider}, connections: map[string]*runtimeConnection{}}
lifecycle := &runtimeLifecycle{store: writer, lease: leased.Lease, registry: registry, config: RuntimeProvider{InstallationID: installation, Provider: provider}, connections: map[string]*runtimeConnection{}}
if checkpoint {
lifecycle.config.Provider = waitingCleanupCheckpoint{beforeKill: provider.beforeKill}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ import (
type finishObservationFixture struct {
s *store.Store
writer *store.Store
lease Ownership
pool *pgxpool.Pool
tenant string
session store.Session
Expand All @@ -28,7 +29,7 @@ type finishObservationFixture struct {
func newFinishObservationFixture(t *testing.T, maxConnections int32) finishObservationFixture {
t.Helper()
var cfg *pgxpool.Config
s, writer := resetManagerStoreConfig(t, func(c *pgxpool.Config) {
s, owner := resetManagerStoreConfig(t, func(c *pgxpool.Config) {
if maxConnections > 0 {
c.MaxConns = maxConnections
}
Expand Down Expand Up @@ -56,7 +57,7 @@ func newFinishObservationFixture(t *testing.T, maxConnections int32) finishObser
if err != nil {
t.Fatal(err)
}
return finishObservationFixture{s, writer, pool, tenant, session, Dispatcher{Store: writer}}
return finishObservationFixture{s, owner.Store, owner.Lease, pool, tenant, session, Dispatcher{Store: owner.Store}}
}
func (f finishObservationFixture) start(t *testing.T) store.InputReceipt {
t.Helper()
Expand Down Expand Up @@ -183,7 +184,7 @@ func TestFinishRunObservationLockTimeoutAndFailureKeepLease(t *testing.T) {
if elapsed := time.Since(started); elapsed < 900*time.Millisecond || elapsed > 2*time.Second {
t.Fatal("pool timeout exceeded observation budget", elapsed)
}
if err = f.writer.CheckExecutionOwnership(t.Context()); err != nil {
if err = f.lease.CheckOwnership(t.Context()); err != nil {
t.Fatal("observation cancelled lease", err)
}
return
Expand Down Expand Up @@ -211,7 +212,7 @@ func TestFinishRunObservationLockTimeoutAndFailureKeepLease(t *testing.T) {
if used != nil || code != nil {
t.Fatal("failed observation wrote")
}
if err = f.writer.CheckExecutionOwnership(t.Context()); err != nil {
if err = f.lease.CheckOwnership(t.Context()); err != nil {
t.Fatal("observation cancelled execution lease", err)
}
})
Expand Down
33 changes: 33 additions & 0 deletions services/core/internal/execution/owner.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
package execution

import (
"context"
"time"

"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/store"
)

// Ownership is the execution lease as the Worker uses it. *pgunit.Lease implements it.
type Ownership interface {
CheckOwnership(context.Context) error
CancelOperations(context.Context, context.CancelFunc) error
Close(context.Context) error
}

// Owner is everything bound to one execution lease. Later cutovers add one explicit
// field per domain's execution operations and delete the matching store calls.
type Owner struct {
Lease Ownership
Store *store.Store // the remaining store execution operations, built by store.NewExecution(s, lease)
}

// leaseCloseTimeout bounds releasing the lease once the Worker owns it.
const leaseCloseTimeout = 5 * time.Second

// closeLease releases the lease within its own deadline, independent of the
// cancellation of the request or run that ends ownership.
func closeLease(ctx context.Context, lease Ownership) error {
ctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), leaseCloseTimeout)
defer cancel()
return lease.Close(ctx)
}
136 changes: 136 additions & 0 deletions services/core/internal/execution/owner_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,136 @@
package execution

import (
"context"
"errors"
"sync/atomic"
"testing"
"time"

"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway"
"github.com/google/uuid"
)

// closeCountingLease counts Close calls on the lease a Worker owns. It forwards
// to inner; without inner, any call other than Close fails the test.
type closeCountingLease struct {
t *testing.T
inner Ownership
closable atomic.Bool
closes atomic.Int32
}

func (l *closeCountingLease) CheckOwnership(ctx context.Context) error {
if l.inner == nil {
l.t.Error("unexpected call to CheckOwnership")
return errors.New("unexpected call to CheckOwnership")
}
return l.inner.CheckOwnership(ctx)
}

func (l *closeCountingLease) CancelOperations(ctx context.Context, cancel context.CancelFunc) error {
if l.inner == nil {
l.t.Error("unexpected call to CancelOperations")
return errors.New("unexpected call to CancelOperations")
}
return l.inner.CancelOperations(ctx, cancel)
}

func (l *closeCountingLease) Close(ctx context.Context) error {
l.closes.Add(1)
if !l.closable.Load() {
l.t.Error("lease closed before its owner finished")
}
if _, bounded := ctx.Deadline(); !bounded || ctx.Err() != nil {
l.t.Error("lease closed without a live bounded context", ctx.Err())
}
if l.inner == nil {
return nil
}
return l.inner.Close(ctx)
}

func TestStartWorkerFailureClosesLeaseOnce(t *testing.T) {
if _, err := StartWorker(t.Context(), &Dispatcher{}, Owner{}); err == nil {
t.Fatal("worker started without an execution lease")
}
// The failed request's context is already canceled; the close must not inherit it.
canceled, cancel := context.WithCancel(t.Context())
cancel()
for name, start := range map[string]func(*testing.T, *closeCountingLease) error{
"negative concurrency": func(t *testing.T, lease *closeCountingLease) error {
_, err := StartWorker(canceled, &Dispatcher{MaxConcurrentExecutions: -1}, Owner{Lease: lease})
return err
},
"excess concurrency": func(t *testing.T, lease *closeCountingLease) error {
_, err := StartWorker(canceled, &Dispatcher{MaxConcurrentExecutions: 1025}, Owner{Lease: lease})
return err
},
"missing Store": func(t *testing.T, lease *closeCountingLease) error {
_, err := StartWorker(canceled, &Dispatcher{}, Owner{Lease: lease})
return err
},
"deployment claim": func(t *testing.T, lease *closeCountingLease) error {
s, owner := resetManagerStore(t)
lease.inner = owner.Lease
id := uuid.NewString()
dispatcher := &Dispatcher{Store: s, Registry: runtimegateway.NewRegistry(), ManagedRuntimes: NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil })}
_, err := StartWorker(canceled, dispatcher, Owner{Lease: lease, Store: owner.Store})
if ping := owner.Lease.CheckOwnership(t.Context()); !errors.Is(ping, pgunit.ErrLeaseClosed) {
t.Error("failed start kept the database lease", ping)
}
return err
},
} {
t.Run(name, func(t *testing.T) {
lease := &closeCountingLease{t: t}
lease.closable.Store(true)
if err := start(t, lease); err == nil {
t.Fatal("worker started")
}
if closes := lease.closes.Load(); closes != 1 {
t.Fatal("failed start closed the lease", closes, "times")
}
})
}
}

func TestWorkerRunClosesLeaseAfterDrain(t *testing.T) {
s, owner := resetManagerStore(t)
lease := &closeCountingLease{t: t, inner: owner.Lease}
id := uuid.NewString()
dispatcher := &Dispatcher{Store: s, Registry: runtimegateway.NewRegistry(), ManagedRuntimes: NewDeferredRuntimeProvider(id, func(context.Context) (*RuntimeProvider, error) { return nil, nil })}
worker, err := StartWorker(t.Context(), dispatcher, Owner{Lease: lease, Store: owner.Store})
if err != nil {
t.Fatal(err)
}
// An external provisioning caller is still in flight when Run exits.
worker.runtimes.active.Add(1)
ctx, cancel := context.WithCancel(t.Context())
done := make(chan error, 1)
go func() { done <- worker.Run(ctx) }()
cancel()
select {
case <-worker.runtimes.ctx.Done():
case <-time.After(5 * time.Second):
worker.runtimes.active.Done()
t.Fatal("Run did not stop its runtimes")
}
lease.closable.Store(true)
worker.runtimes.active.Done()
select {
case err := <-done:
if !errors.Is(err, context.Canceled) {
t.Fatal(err)
}
case <-time.After(5 * time.Second):
t.Fatal("Run did not exit after draining")
}
if closes := lease.closes.Load(); closes != 1 {
t.Fatal("Run closed the lease", closes, "times")
}
if ping := owner.Lease.CheckOwnership(t.Context()); !errors.Is(ping, pgunit.ErrLeaseClosed) {
t.Fatal("Run kept the database lease", ping)
}
}
7 changes: 4 additions & 3 deletions services/core/internal/execution/prepared_dispatch.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,9 +16,10 @@ type EnvironmentRun struct {
Turn store.Turn
}

// RunEnvironmentInput reserves a Turn on the Session-owned Runtime Executor.
func (d *Dispatcher) RunEnvironmentInput(ctx context.Context, tenantID, sessionID, reservationID string) (run EnvironmentRun, err error) {
if err = d.Store.CheckExecutionOwnership(ctx); err != nil {
// RunEnvironmentInput reserves a Turn on the Session-owned Runtime Executor. It
// checks lease, the lease d.Store was built on, before any Runtime preparation.
func (d *Dispatcher) RunEnvironmentInput(ctx context.Context, lease Ownership, tenantID, sessionID, reservationID string) (run EnvironmentRun, err error) {
if err = lease.CheckOwnership(ctx); err != nil {
return run, err
}
run.Reservation, err = d.Store.ExpireEnvironmentInput(ctx, tenantID, sessionID, reservationID)
Expand Down
7 changes: 1 addition & 6 deletions services/core/internal/execution/runtime_cancellation.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,12 +14,7 @@ func (m *runtimeManager) cancelLifecycles(nodes []*runtimeNode) error {
n.lifecycle.cancelOperations()
}
}
// Managers in pure lifecycle tests have no Store or execution connection.
if m.store == nil {
cancel()
return nil
}
err := m.store.CancelExecutionOperations(m.ctx, cancel)
err := m.lease.CancelOperations(m.ctx, cancel)
if err != nil {
// A manual reconcile may have no running coordinator to consume failed.
// Close admission synchronously; Worker shutdown still owns cancellation.
Expand Down
2 changes: 1 addition & 1 deletion services/core/internal/execution/runtime_compute.go
Original file line number Diff line number Diff line change
Expand Up @@ -94,7 +94,7 @@ func (r *runtimeLifecycle) observeCompute(ctx context.Context, owner store.Runti
if owner.SessionDeleted || owner.Expired || owner.State == "cleanup_pending" {
return r.cleanupCompute(ctx, p, owner, state)
}
if err := r.store.CheckExecutionOwnership(ctx); err != nil {
if err := r.lease.CheckOwnership(ctx); err != nil {
return err
}
switch owner.ComputePhase {
Expand Down
2 changes: 1 addition & 1 deletion services/core/internal/execution/runtime_compute_wake.go
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,7 @@ func (r *runtimeLifecycle) wakeCompute(ctx context.Context, p sandbox.Checkpoint
}

func (r *runtimeLifecycle) cleanupCompute(ctx context.Context, p sandbox.CheckpointProvider, owner store.RuntimeAllocation, state runtimeCompute) error {
if err := r.store.CheckExecutionOwnership(ctx); err != nil {
if err := r.lease.CheckOwnership(ctx); err != nil {
return err
}
// An uncommitted artifact is found by its persisted attempt, never a directory
Expand Down
Loading
Loading