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
37 changes: 31 additions & 6 deletions apps/daemon/internal/cli/connect_suspend_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -339,7 +339,7 @@ func TestSuspensionReconnectBeforeConfirmation(t *testing.T) {
}
}

func TestQuiesceCloseFailureDisconnectsAndSettlesBeforeReconnect(t *testing.T) {
func TestQuiesceRetainsExecutorAcrossReconnectAndShutdownConfirmsCleanup(t *testing.T) {
environment, session := uuid.NewString(), uuid.NewString()
t.Setenv("OAC_RUNTIME_HOME", t.TempDir())
t.Setenv("OAC_RUNTIME_STATE_DIRECTORY", t.TempDir())
Expand Down Expand Up @@ -431,13 +431,38 @@ func TestQuiesceCloseFailureDisconnectsAndSettlesBeforeReconnect(t *testing.T) {
}
readStatus("released")
sendLifecycleFrame(t, peer, proto.TypeEnvironmentQuiesce, request)
if got := readLifecycleResult(t, peer, proto.TypeEnvironmentQuiesced); got.Accepted || got.ErrorCode != "resource_busy" {
t.Fatalf("failed Close acknowledged: %+v", got)
if got := readLifecycleResult(t, peer, proto.TypeEnvironmentQuiesced); !got.Accepted {
t.Fatalf("idle owner rejected: %+v", got)
}
if native.closes.Load() != 0 {
t.Fatal("quiesce closed the native owner")
}
control.signal <- syscall.SIGUSR1
select {
case peer = <-peers:
case <-time.After(3 * time.Second):
t.Fatal("no planned reconnect")
}
defer peer.Close()
sendLifecycleFrame(t, peer, proto.TypeEnvironmentResume, request)
if got := readLifecycleResult(t, peer, proto.TypeEnvironmentResumed); !got.Accepted {
t.Fatalf("resume rejected: %+v", got)
}
preparation.ID = "resumed-prepare"
if err := peer.WriteJSON(preparation); err != nil {
t.Fatal(err)
}
resumed := readStatus("ready")
if !resumed.Reused || resumed.ExecutorID != ready.ExecutorID || resumed.Handle == ready.Handle || native.closes.Load() != 0 {
t.Fatalf("resume replaced native owner: before=%+v after=%+v closes=%d", ready, resumed, native.closes.Load())
}
// Cancelling the daemon still retires resident native resources. Failed
// cleanup must settle before the loop can finish or create another router.
cancel()
select {
case <-native.retry:
case <-time.After(3 * time.Second):
t.Fatal("fenced close failure did not enter Shutdown retry")
t.Fatal("cancelled resident owner did not enter Shutdown retry")
}
_ = peer.SetReadDeadline(time.Now().Add(time.Second))
for {
Expand All @@ -462,7 +487,7 @@ func TestQuiesceCloseFailureDisconnectsAndSettlesBeforeReconnect(t *testing.T) {
if native.closes.Load() != 2 {
t.Fatalf("Close calls=%d", native.closes.Load())
}
if _, err := os.Stat(control.path); !os.IsNotExist(err) {
t.Fatal("failed quiesce armed snapshot control")
if control.identity.SuspendID != "" || control.lastResumed == nil || !control.lastResumed.SameSuspension(request) {
t.Fatal("resumed suspension left snapshot control armed or lost its receipt")
}
}
4 changes: 2 additions & 2 deletions apps/daemon/internal/dispatch/executor.go
Original file line number Diff line number Diff line change
Expand Up @@ -253,7 +253,7 @@ func (r *Router) abandonExecutorAdmission(p *preparationState, state, code strin
}

func (r *Router) scheduleExecutorIdleLocked(owner *executorState) {
if owner.invalid || owner.preparing || owner.run != nil || owner.admission != nil {
if r.closed || r.suspension != nil || owner.invalid || owner.preparing || owner.run != nil || owner.admission != nil {
return
}
if owner.timer != nil {
Expand Down Expand Up @@ -281,7 +281,7 @@ func (r *Router) scheduleExecutorIdleLocked(owner *executorState) {
// expireIdleExecutor closes owner only while lease is its current idle lease.
func (r *Router) expireIdleExecutor(owner *executorState, lease uint64) {
r.mu.Lock()
if r.executors[owner.sessionID] != owner || owner.idleLease != lease || owner.run != nil || owner.admission != nil || owner.invalid {
if r.closed || r.suspension != nil || r.executors[owner.sessionID] != owner || owner.idleLease != lease || owner.run != nil || owner.admission != nil || owner.invalid {
r.mu.Unlock()
return
}
Expand Down
49 changes: 49 additions & 0 deletions apps/daemon/internal/dispatch/executor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -342,3 +342,52 @@ func TestExecutorRejectsOutputFromAnotherTurn(t *testing.T) {
func (*reusableTurn) SteerWithReceipt(context.Context, proto.PromptSteerPayload, func()) error {
return agent.ErrSteeringRejected
}

func TestExecutorResumesSameOwnerAndNextTurnAfterSuspension(t *testing.T) {
owner := &reusableExecutor{starts: make(chan *reusableTurn, 2)}
var factories atomic.Int32
registry := agent.NewRegistry()
registerExecutorKind(registry, proto.SupportedAgentKind{Kind: "prepared", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{LocalEnvironment: proto.CapabilitySupported, DurableInputReceipts: proto.CapabilitySupported})}, func(context.Context, proto.PromptRequestPayload) (agent.Executor, error) {
factories.Add(1)
return owner, nil
})
sender := &recSender{}
r, err := dispatch.New(dispatch.Config{Registry: registry, Sender: sender, LocalWorkspace: preparationWorkspace(t)})
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
if err := r.Shutdown(context.Background()); err != nil {
t.Error(err)
}
})
first := executorAdmission(t, r, sender, "first", preparationRequest())
startExecutorTurn(t, r, sender, "first", "one", first)
(<-owner.starts).finish()
waitFor(t, func() bool { return r.ActiveRuns() == 0 }, "first Turn settled")
request := proto.EnvironmentSuspendPayload{EnvironmentID: preparationEnvironmentID, SuspendID: "suspension"}
if err := r.Quiesce(t.Context(), request); err != nil {
t.Fatal(err)
}
resumedSender := &recSender{}
if err := r.Resume(request, resumedSender); err != nil {
t.Fatal(err)
}
req := preparationRequest()
req.Configuration.AgentSessionID = "native-session"
conflict := req
conflict.Configuration.AgentOptions = map[string]any{"model": "changed"}
if err := r.Handle(t.Context(), mustEnv(t, proto.TypeExecutionPrepare, "conflict", conflict)); err == nil {
t.Fatal("suspension relaxed the fixed configuration")
}
second := executorAdmission(t, r, resumedSender, "second", req)
if !second.Reused || first.ExecutorID != second.ExecutorID || first.Handle == second.Handle || factories.Load() != 1 || owner.closes.Load() != 0 {
t.Fatal("suspension replaced the resident native owner")
}
startExecutorTurn(t, r, resumedSender, "second", "two", second)
(<-owner.starts).finish()
waitFor(t, func() bool { return r.ActiveRuns() == 0 }, "resumed Turn settled")
if owner.closes.Load() != 0 {
t.Fatal("resumed Turn retired its healthy native owner")
}
}
26 changes: 9 additions & 17 deletions apps/daemon/internal/dispatch/suspend.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,6 @@ package dispatch
import (
"context"
"errors"
"fmt"
"strings"

"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto"
Expand All @@ -14,9 +13,9 @@ var ErrRouterQuiesced = errors.New("dispatch: router quiesced")
var ErrRouterBusy = errors.New("dispatch: router has unsettled work")

// Quiesce serializes against admission, then drains every admitted output and
// receipt and closes idle Executors before acknowledging suspension. Busy
// rejection leaves admission open; a failed close or drain timeout keeps it
// closed until the caller shuts the connection down.
// receipt while retaining settled idle Executors for suspension. Busy rejection
// leaves admission open; a drain timeout keeps it closed until the caller shuts
// the connection down.
func (r *Router) Quiesce(ctx context.Context, request proto.EnvironmentSuspendPayload) error {
if strings.TrimSpace(request.EnvironmentID) == "" || strings.TrimSpace(request.SuspendID) == "" || len(request.SuspendID) > 128 {
return errors.New("dispatch: invalid suspension identity")
Expand Down Expand Up @@ -56,27 +55,20 @@ func (r *Router) Quiesce(ctx context.Context, request proto.EnvironmentSuspendPa
p.timer.Stop()
}
}
owners := r.closeIdleExecutorsLocked()
for _, owner := range owners {
for _, owner := range r.executors {
if owner.timer != nil {
owner.timer.Stop()
}
// Stop cannot retract a callback already waiting for the lock. Fence its
// lease before releasing the lock, including after an eventual Resume.
owner.idleLease++
owner.closeReason = "suspend"
}
r.mu.Unlock()
r.closeIdleExecutors(owners)
err := r.shutdownWG.waitContext(ctx)
r.mu.Lock()
if r.closed {
err = ErrRouterClosed
} else if err == nil {
for _, owner := range r.executors {
cause := owner.closeErr
if cause == nil {
cause = errors.New("cleanup has not settled")
}
err = errors.Join(err, fmt.Errorf("dispatch: executor %s: %w", owner.id, cause))
}
}

r.mu.Unlock()
if err != nil {
return errors.Join(ErrRouterQuiesced, err)
Expand Down
Loading
Loading