From 5546601f9adbc47defe9c2e53563475c8e603424 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Thu, 8 Oct 2026 07:56:11 +0000 Subject: [PATCH] Settle the agent-host placement review findings An idle pass consults the assignment only to quiesce, so an unplaced Session or a disconnected agent host is not recorded as a sandbox fault. The initialization scanner places a Session only on a host that advertises its Harness as available, as the Turn path does. waitServing returns a shutdown's own error and leaves the pass time to resume. TestInitializationBindsAgentHost uses a configured deployment, which hosted admission now requires. --- .../internal/execution/runtime_compute.go | 26 +++++++++++++------ .../execution/runtime_compute_wake.go | 13 +++++++--- .../execution/runtime_initialization.go | 6 +++-- .../tests/integration/link_authority_test.go | 3 ++- .../runtime_compute_lifecycle_test.go | 5 ++++ 5 files changed, 38 insertions(+), 15 deletions(-) diff --git a/services/core/internal/execution/runtime_compute.go b/services/core/internal/execution/runtime_compute.go index adab45df7..db126771f 100644 --- a/services/core/internal/execution/runtime_compute.go +++ b/services/core/internal/execution/runtime_compute.go @@ -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" @@ -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 } @@ -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) diff --git a/services/core/internal/execution/runtime_compute_wake.go b/services/core/internal/execution/runtime_compute_wake.go index 8458f240c..c7aacd58f 100644 --- a/services/core/internal/execution/runtime_compute_wake.go +++ b/services/core/internal/execution/runtime_compute_wake.go @@ -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: } diff --git a/services/core/internal/execution/runtime_initialization.go b/services/core/internal/execution/runtime_initialization.go index 531592e90..8d836c44d 100644 --- a/services/core/internal/execution/runtime_initialization.go +++ b/services/core/internal/execution/runtime_initialization.go @@ -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 } diff --git a/services/core/tests/integration/link_authority_test.go b/services/core/tests/integration/link_authority_test.go index 2488a5edf..7ada0a95d 100644 --- a/services/core/tests/integration/link_authority_test.go +++ b/services/core/tests/integration/link_authority_test.go @@ -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 + `}`), diff --git a/services/core/tests/integration/runtime_compute_lifecycle_test.go b/services/core/tests/integration/runtime_compute_lifecycle_test.go index 762b29919..d595b9360 100644 --- a/services/core/tests/integration/runtime_compute_lifecycle_test.go +++ b/services/core/tests/integration/runtime_compute_lifecycle_test.go @@ -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 {