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
26 changes: 18 additions & 8 deletions services/core/internal/execution/runtime_compute.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import (
"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/deployment"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/providercontract"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions"
"github.com/google/uuid"
Expand Down Expand Up @@ -119,14 +120,6 @@ func (r *runtimeLifecycle) idleCompute(ctx context.Context, p sandbox.SandboxPro
if compute.Status != "running" || !compute.BootstrapComplete {
return sandbox.ErrComputeUnconfirmed
}
bound, err := r.allocationAssignment(ctx, owner)
if err != nil {
return err
}
peer, err := assignedRuntimePeer(ctx, r.sessions, r.registry, bound)
if err != nil {
return err
}
if _, err := r.deployment.CheckRunning(ctx, owner); err != nil {
return err
}
Expand All @@ -142,6 +135,23 @@ func (r *runtimeLifecycle) idleCompute(ctx context.Context, p sandbox.SandboxPro
if policy == nil || !activity.ReadyToSuspend(policy.IdleTimeout) {
return nil
}
// Only the assigned agent host quiesces the Environment. An unplaced
// Session has nothing to quiesce, and a later pass retries a host that is
// not connected; neither is a sandbox fault.
bound, err := r.allocationAssignment(ctx, owner)
if errors.Is(err, sessions.ErrNotFound) {
return nil
}
if err != nil {
return err
}
peer, err := assignedRuntimePeer(ctx, r.sessions, r.registry, bound)
if errors.Is(err, sessions.ErrNotFound) || errors.Is(err, runtimegateway.ErrSessionClosed) || errors.Is(err, runtimegateway.ErrDeviceNotRegistered) {
return nil
}
if err != nil {
return err
}
state.SuspendID, state.RestoreID, state.Rollback = uuid.NewString(), "", false
until := activity.ObservedAt.Add(policy.Retention)
next, err := r.saveCompute(ctx, owner, "quiescing", state, &until)
Expand Down
13 changes: 9 additions & 4 deletions services/core/internal/execution/runtime_compute_wake.go
Original file line number Diff line number Diff line change
Expand Up @@ -55,16 +55,21 @@ func (r *runtimeLifecycle) wakeCompute(ctx context.Context, p sandbox.SandboxPro
return r.deployment.ClearWake(ctx, next, owner.ComputeActivityAt)
}

// waitServing waits, for at most 30 seconds, until the relay holds the serve
// peer of the allocation's Link resource.
// waitServing waits until the relay holds the serve peer of the allocation's
// Link resource, for at most 20 seconds, which leaves the pass's 30-second
// operation time to resume the Environment. A shutdown or lost lease returns
// its own error rather than a compute diagnostic.
func (r *runtimeLifecycle) waitServing(ctx context.Context, owner deployment.Allocation) error {
ctx, cancel := context.WithTimeout(ctx, 30*time.Second)
wait, cancel := context.WithTimeout(ctx, 20*time.Second)
defer cancel()
timer := time.NewTicker(100 * time.Millisecond)
defer timer.Stop()
for !r.links.Serving(serveResource(owner).Ref()) {
select {
case <-ctx.Done():
case <-wait.Done():
if errors.Is(ctx.Err(), context.Canceled) {
return ctx.Err()
}
return sandbox.ErrComputeUnconfirmed
case <-timer.C:
}
Expand Down
6 changes: 4 additions & 2 deletions services/core/internal/execution/runtime_initialization.go
Original file line number Diff line number Diff line change
Expand Up @@ -67,12 +67,14 @@ func (w *Worker) runEnvironmentInitializations(ctx context.Context) error {
// Placement binds an agent host once the Environment serves;
// a later scan claims the preparation on it.
if _, err := w.place(ctx, owner.TenantID, owner.SessionID, owner.EnvironmentID, func(id string) bool {
// The Turn path's predicate: the host advertises the
// Harness as available.
peer, err := w.dispatcher.authorizedPeer(ctx, id)
if err != nil {
return false
}
_, _, known := peer.AgentKindStatus(owner.Engine)
return known
_, err = runtimeDeclaration(peer, owner.Engine)
return err == nil
}); err != nil && !errors.Is(err, sessions.ErrNotFound) && !errors.Is(err, sessions.ErrDeviceBindingConflict) {
return err
}
Expand Down
3 changes: 2 additions & 1 deletion services/core/tests/integration/link_authority_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -536,7 +536,8 @@ func TestRegisteredAgentHostAuthenticates(t *testing.T) {
func TestInitializationBindsAgentHost(t *testing.T) {
for _, environment := range []string{`{"type":"openai_hosted"}`, `{"type":"self_hosted","workspace_directory":"/workspace"}`} {
t.Run(environment, func(t *testing.T) {
s, _ := newManagedTestStore(t)
// Hosted work is admitted only on a configured deployment.
s, _ := configuredStore(t)
tenant, host := uuid.NewString(), registerAgentHost(t, s, "")
session, err := s.CreateSession(t.Context(), tenant, WithFixtureModelProvider(sessions.CreateSession{Creator: FixtureCreator(), Engine: "codex", IdempotencyKey: uuid.NewString(),
Configuration: json.RawMessage(`{"agent":{"model":"test-model"},"environment":` + environment + `}`),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -440,6 +440,11 @@ func TestRuntimeComputeLifecycleIdleSuspendAndQueuedSameSessionWake(t *testing.T
if f.provider.captures != 0 {
t.Fatal("never-used Session suspended")
}
// Before placement there is no agent host to ask, which is not a sandbox fault.
var diagnostic string
if err := f.store.pool.QueryRow(t.Context(), `SELECT observation_error FROM runtime_allocations WHERE id=$1`, owner.ID).Scan(&diagnostic); err != nil || diagnostic != "" {
t.Fatalf("unplaced Session observed as %q: %v", diagnostic, err)
}
completed := f.complete(owner)
suspended := f.phase(tenant, env.ID, "suspended")
if f.provider.captures != 1 || f.provider.computeKills != 1 || len(f.provider.computes) != 0 || len(f.provider.snapshots) != 1 {
Expand Down
Loading