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
20 changes: 14 additions & 6 deletions services/core/internal/execution/runtime_compute.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -201,14 +206,16 @@ 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
}
_, err = r.saveCompute(ctx, next, "suspended", state, next.ComputeRetainedUntil)
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
Expand All @@ -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
}
Expand Down
7 changes: 5 additions & 2 deletions services/core/internal/execution/runtime_compute_wake.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down
25 changes: 25 additions & 0 deletions services/core/internal/execution/runtime_observation.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down
93 changes: 93 additions & 0 deletions services/core/internal/execution/runtime_observation_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
})
}
}
Loading