From 4b9bf0d1fb41c39f889145f52c33abfa7803048c Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Sat, 10 Oct 2026 17:23:54 +0000 Subject: [PATCH] Retain compute transition failure observations --- .../internal/execution/runtime_compute.go | 20 ++-- .../execution/runtime_compute_wake.go | 7 +- .../internal/execution/runtime_observation.go | 25 +++++ .../execution/runtime_observation_test.go | 93 +++++++++++++++++++ 4 files changed, 137 insertions(+), 8 deletions(-) create mode 100644 services/core/internal/execution/runtime_observation_test.go diff --git a/services/core/internal/execution/runtime_compute.go b/services/core/internal/execution/runtime_compute.go index 9512d6765..f558ef275 100644 --- a/services/core/internal/execution/runtime_compute.go +++ b/services/core/internal/execution/runtime_compute.go @@ -45,7 +45,8 @@ func (r *runtimeLifecycle) saveCompute(ctx context.Context, owner deployment.All } return r.deployment.SetCompute(ctx, owner, phase, raw, until, idleTimeout) } -func (r *runtimeLifecycle) enableCompute(ctx context.Context, owner deployment.Allocation) error { +func (r *runtimeLifecycle) enableCompute(ctx context.Context, owner deployment.Allocation) (err error) { + defer func() { err = withObservationOwner(owner, err) }() p := r.config.Provider if err := providercontract.Require(p, "Initial"); err != nil { return err @@ -65,7 +66,8 @@ func (r *runtimeLifecycle) enableCompute(ctx context.Context, owner deployment.A return err } -func (r *runtimeLifecycle) observeCompute(ctx context.Context, owner deployment.Allocation) error { +func (r *runtimeLifecycle) observeCompute(ctx context.Context, owner deployment.Allocation) (err error) { + defer func() { err = withObservationOwner(owner, err) }() p := r.config.Provider if err := providercontract.Require(p, "Initial"); err != nil { return err @@ -105,7 +107,8 @@ func (r *runtimeLifecycle) observeCompute(ctx context.Context, owner deployment. } } -func (r *runtimeLifecycle) idleCompute(ctx context.Context, p sandbox.SandboxProvider, owner deployment.Allocation, state runtimeCompute) error { +func (r *runtimeLifecycle) idleCompute(ctx context.Context, p sandbox.SandboxProvider, owner deployment.Allocation, state runtimeCompute) (err error) { + defer func() { err = withObservationOwner(owner, err) }() compute, err := p.GetCompute(ctx, runtimeReference(owner), state.Current) if err != nil { return err @@ -141,6 +144,7 @@ func (r *runtimeLifecycle) idleCompute(ctx context.Context, p sandbox.SandboxPro if err != nil { return err } + owner = next result, err := peer.SuspendControl(ctx, proto.TypeEnvironmentQuiesce, proto.EnvironmentSuspendPayload{EnvironmentID: owner.EnvironmentID, SuspendID: state.SuspendID}) if err != nil { return err @@ -175,7 +179,8 @@ func (r *runtimeLifecycle) idleCompute(ctx context.Context, p sandbox.SandboxPro return r.captureCompute(ctx, p, suspending, state, false) } -func (r *runtimeLifecycle) captureCompute(ctx context.Context, p sandbox.SandboxProvider, owner deployment.Allocation, state runtimeCompute, observeOnly bool) error { +func (r *runtimeLifecycle) captureCompute(ctx context.Context, p sandbox.SandboxProvider, owner deployment.Allocation, state runtimeCompute, observeOnly bool) (err error) { + defer func() { err = withObservationOwner(owner, err) }() result, err := p.Suspend(ctx, sandbox.SuspendRequest{Reference: runtimeReference(owner), OperationID: state.SuspendID, Source: state.Current, Snapshot: state.Snapshot, ObserveOnly: observeOnly}) if err != nil { return err @@ -201,6 +206,7 @@ func (r *runtimeLifecycle) captureCompute(ctx context.Context, p sandbox.Sandbox if err != nil { return err } + owner = next if err := ignoreComputeAbsent(p.KillCompute(ctx, runtimeReference(owner), state.Current)); err != nil { return err } @@ -208,7 +214,8 @@ func (r *runtimeLifecycle) captureCompute(ctx context.Context, p sandbox.Sandbox return err } -func (r *runtimeLifecycle) restoreIdleCompute(ctx context.Context, p sandbox.SandboxProvider, owner deployment.Allocation, state runtimeCompute) error { +func (r *runtimeLifecycle) restoreIdleCompute(ctx context.Context, p sandbox.SandboxProvider, owner deployment.Allocation, state runtimeCompute) (err error) { + defer func() { err = withObservationOwner(owner, err) }() activity, err := r.reader.Activity(ctx, owner.ID) if err != nil { return err @@ -231,7 +238,8 @@ func (r *runtimeLifecycle) restoreIdleCompute(ctx context.Context, p sandbox.San } return r.restoreCompute(ctx, p, next, state, false) } -func (r *runtimeLifecycle) restoreCompute(ctx context.Context, p sandbox.SandboxProvider, owner deployment.Allocation, state runtimeCompute, observeOnly bool) error { +func (r *runtimeLifecycle) restoreCompute(ctx context.Context, p sandbox.SandboxProvider, owner deployment.Allocation, state runtimeCompute, observeOnly bool) (err error) { + defer func() { err = withObservationOwner(owner, err) }() if state.Target == nil || state.Snapshot == nil || state.Rollback { return sandbox.ErrOwnership } diff --git a/services/core/internal/execution/runtime_compute_wake.go b/services/core/internal/execution/runtime_compute_wake.go index 315b08c37..2aab604f6 100644 --- a/services/core/internal/execution/runtime_compute_wake.go +++ b/services/core/internal/execution/runtime_compute_wake.go @@ -13,7 +13,8 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" ) -func (r *runtimeLifecycle) wakeCompute(ctx context.Context, p sandbox.SandboxProvider, owner deployment.Allocation, state runtimeCompute) error { +func (r *runtimeLifecycle) wakeCompute(ctx context.Context, p sandbox.SandboxProvider, owner deployment.Allocation, state runtimeCompute) (err error) { + defer func() { err = withObservationOwner(owner, err) }() if state.Rollback { if _, err := p.ResumeCompute(ctx, runtimeReference(owner), state.Current); err != nil { return err @@ -65,7 +66,9 @@ func (r *runtimeLifecycle) wakeCompute(ctx context.Context, p sandbox.SandboxPro if err != nil { return err } - if err := r.deployment.ClearWake(ctx, next, owner.ComputeActivityAt); err != nil { + observedActivity := owner.ComputeActivityAt + owner = next + if err := r.deployment.ClearWake(ctx, next, observedActivity); err != nil { return err } return r.observeConnection(ctx, next) diff --git a/services/core/internal/execution/runtime_observation.go b/services/core/internal/execution/runtime_observation.go index e4d93e3d0..bc8cf2fab 100644 --- a/services/core/internal/execution/runtime_observation.go +++ b/services/core/internal/execution/runtime_observation.go @@ -8,7 +8,32 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" ) +// observationFailure retains the receipt held by the operation that failed. +// A later phase must never be queried and relabelled with an earlier failure. +type observationFailure struct { + owner deployment.Allocation + cause error +} + +func (e *observationFailure) Error() string { return e.cause.Error() } +func (e *observationFailure) Unwrap() error { return e.cause } + +func withObservationOwner(owner deployment.Allocation, err error) error { + if err == nil { + return nil + } + var observed *observationFailure + if errors.As(err, &observed) { + return err + } + return &observationFailure{owner: owner, cause: err} +} + func (r *runtimeLifecycle) recordObservation(ctx context.Context, owner deployment.Allocation, observed error) { + var failure *observationFailure + if errors.As(observed, &failure) { + owner = failure.owner + } if owner.NodeID == "" { return } diff --git a/services/core/internal/execution/runtime_observation_test.go b/services/core/internal/execution/runtime_observation_test.go new file mode 100644 index 000000000..8ad10bdae --- /dev/null +++ b/services/core/internal/execution/runtime_observation_test.go @@ -0,0 +1,93 @@ +package execution + +import ( + "context" + "errors" + "testing" + "time" + + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" +) + +type observationCaptureProvider struct { + retentionProvider + captureError error +} + +func (p *observationCaptureProvider) Suspend(_ context.Context, q sandbox.SuspendRequest) (sandbox.ComputeState, error) { + if p.captureError != nil { + return sandbox.ComputeState{}, p.captureError + } + return sandbox.ComputeState{Compute: q.Source, Snapshot: &sandbox.SnapshotIdentity{ID: "captured"}, Status: "suspended", SourceStopped: true}, nil +} + +func TestComputeFailureObservationUsesCommittedReceipt(t *testing.T) { + for _, stage := range []string{"capture", "source_cleanup"} { + t.Run(stage, func(t *testing.T) { + provider := &observationCaptureProvider{} + if stage == "capture" { + provider.captureError = sandbox.ErrComputeUnconfirmed + } else { + provider.killError = sandbox.ErrComputeUnconfirmed + } + fixture := workspaceSettlementFixtureWithSetup(t, provider, workspaceNodeSetup(t, true)) + r, session := fixture.lifecycle, fixture.session + initial, err := r.provision(t.Context(), session.TenantID, session.Environment.ID, r.config.InstallationID) + if err != nil { + t.Fatal(err) + } + // Seed the acknowledged idle barrier; subsequent phase commits use the real + // Session-locked lifecycle operations, with real observation CAS writes. + if _, err = fixture.pool.Exec(t.Context(), `UPDATE environments SET initialization='complete',status='disconnected' WHERE id=$1`, initial.EnvironmentID); err != nil { + t.Fatal(err) + } + if _, err = fixture.pool.Exec(t.Context(), `UPDATE runtime_allocations SET compute_phase='quiescing',compute_revision=2,compute_retained_until=clock_timestamp()+interval '1 hour' WHERE id=$1`, initial.ID); err != nil { + t.Fatal(err) + } + original, err := r.reader.EnvironmentAllocation(t.Context(), initial.Key()) + if err != nil { + t.Fatal(err) + } + state := runtimeCompute{Current: sandbox.Compute{ID: "source", Name: "source"}, SuspendID: "attempt"} + until := time.Now().Add(time.Hour) + suspending, err := r.saveCompute(t.Context(), original, "suspending", state, &until) + if err != nil { + t.Fatal(err) + } + failure := r.observeCompute(t.Context(), suspending) + if !errors.Is(failure, sandbox.ErrComputeUnconfirmed) { + t.Fatalf("capture failure = %v", failure) + } + // Model the outer reconcile entry retaining its earlier scan receipt. + r.recordObservation(t.Context(), original, failure) + current, err := r.reader.EnvironmentAllocation(t.Context(), initial.Key()) + if err != nil || current.ObservationError != "compute_unconfirmed" { + t.Fatalf("observation after transition = %q, %v", current.ObservationError, err) + } + wantRevision := int64(3) + if stage == "source_cleanup" { + wantRevision++ + } + if current.ComputeRevision != wantRevision { + t.Fatalf("revision = %d, want %d", current.ComputeRevision, wantRevision) + } + // A late result must not borrow a later operation's receipt or erase its + // diagnostic, even when both observations belong to the same allocation. + later, err := r.saveCompute(t.Context(), current, "waking", state, &until) + if err != nil { + t.Fatal(err) + } + r.recordObservation(t.Context(), later, sandbox.ErrOwnership) + r.recordObservation(t.Context(), original, failure) + current, err = r.reader.EnvironmentAllocation(t.Context(), initial.Key()) + if err != nil || current.ObservationError != "ownership_mismatch" { + t.Fatalf("late failure overwrote successor: %q, %v", current.ObservationError, err) + } + r.recordObservation(t.Context(), current, nil) + current, err = r.reader.EnvironmentAllocation(t.Context(), initial.Key()) + if err != nil || current.ObservationError != "" { + t.Fatalf("current successful observation did not clear diagnostic: %q, %v", current.ObservationError, err) + } + }) + } +}