From b63c3b213ced726590b8de638bfa60dc69c80c78 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Sat, 10 Oct 2026 17:39:32 +0000 Subject: [PATCH] Preserve resident executors during suspension --- .../internal/cli/connect_suspend_test.go | 37 +++- apps/daemon/internal/dispatch/executor.go | 4 +- .../daemon/internal/dispatch/executor_test.go | 49 +++++ apps/daemon/internal/dispatch/suspend.go | 26 +-- apps/daemon/internal/dispatch/suspend_test.go | 174 ++++++++++++------ contracts/agents-api/harness-onboarding.md | 2 +- contracts/agents-api/zh/harness-onboarding.md | 4 +- docs/runtime-protocol.md | 2 +- docs/zh/runtime-protocol.md | 4 +- internal/agentdaemon/proto/envelope_test.go | 1 + internal/agentdaemon/proto/suspend.go | 5 +- internal/agentdaemon/proto/version.go | 2 +- 12 files changed, 218 insertions(+), 92 deletions(-) diff --git a/apps/daemon/internal/cli/connect_suspend_test.go b/apps/daemon/internal/cli/connect_suspend_test.go index 52dd2caea..6d92648e2 100644 --- a/apps/daemon/internal/cli/connect_suspend_test.go +++ b/apps/daemon/internal/cli/connect_suspend_test.go @@ -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()) @@ -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 { @@ -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") } } diff --git a/apps/daemon/internal/dispatch/executor.go b/apps/daemon/internal/dispatch/executor.go index 0e47e8c20..3a07e09a2 100644 --- a/apps/daemon/internal/dispatch/executor.go +++ b/apps/daemon/internal/dispatch/executor.go @@ -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 { @@ -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 } diff --git a/apps/daemon/internal/dispatch/executor_test.go b/apps/daemon/internal/dispatch/executor_test.go index a2324a78c..2cba36047 100644 --- a/apps/daemon/internal/dispatch/executor_test.go +++ b/apps/daemon/internal/dispatch/executor_test.go @@ -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") + } +} diff --git a/apps/daemon/internal/dispatch/suspend.go b/apps/daemon/internal/dispatch/suspend.go index 5945c7cdc..b38eb9bdc 100644 --- a/apps/daemon/internal/dispatch/suspend.go +++ b/apps/daemon/internal/dispatch/suspend.go @@ -3,7 +3,6 @@ package dispatch import ( "context" "errors" - "fmt" "strings" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" @@ -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") @@ -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) diff --git a/apps/daemon/internal/dispatch/suspend_test.go b/apps/daemon/internal/dispatch/suspend_test.go index ddb48f030..36562c3e2 100644 --- a/apps/daemon/internal/dispatch/suspend_test.go +++ b/apps/daemon/internal/dispatch/suspend_test.go @@ -123,7 +123,7 @@ func TestQuiesceDrainsPendingReceiptAndFencesConcurrentAdmission(t *testing.T) { } } -func TestQuiesceClosesIdleExecutorOnceAndRequiresExactResume(t *testing.T) { +func TestQuiesceRetainsIdleExecutorAndRequiresExactResume(t *testing.T) { sender := suspendSender(func(context.Context, proto.Envelope) error { return nil }) r := suspensionRouter(t, sender) native := &suspendedExecutor{} @@ -134,32 +134,78 @@ func TestQuiesceClosesIdleExecutorOnceAndRequiresExactResume(t *testing.T) { oldLease := owner.idleLease r.mu.Unlock() request := proto.EnvironmentSuspendPayload{EnvironmentID: "env", SuspendID: "attempt"} - if err := r.Quiesce(context.Background(), request); err != nil { + if err := r.Quiesce(t.Context(), request); err != nil { t.Fatal(err) } r.expireIdleExecutor(owner, oldLease) - if native.closed.Load() != 1 { - t.Fatal("quiesce and stale timer did not close the owner exactly once") + if native.closed.Load() != 0 { + t.Fatal("quiesce or stale timer closed the resident owner") + } + for _, wrong := range []proto.EnvironmentSuspendPayload{ + {EnvironmentID: "env", SuspendID: "obsolete"}, + {EnvironmentID: "foreign", SuspendID: "attempt"}, + } { + if err := r.Resume(wrong, sender); err == nil { + t.Fatal("wrong suspension reopened admission") + } } - wrong := request - wrong.SuspendID = "obsolete" - if err := r.Resume(wrong, sender); err == nil { - t.Fatal("stale operation reopened admission") + if err := r.Resume(request, nil); err == nil { + t.Fatal("resume accepted no sender") } if err := r.Resume(request, sender); err != nil { t.Fatal(err) } + // A callback queued before Stop may run after Resume. It must not close + // the same owner under its renewed idle lease. + r.expireIdleExecutor(owner, oldLease) + r.mu.Lock() + retained := r.executors[owner.sessionID] == owner && !owner.invalid && owner.native == native && owner.idleLease > oldLease + lease := owner.idleLease + r.mu.Unlock() + if !retained || native.closed.Load() != 0 { + t.Fatal("resume replaced or retired the resident owner") + } + r.expireIdleExecutor(owner, lease) + if native.closed.Load() != 1 { + t.Fatal("resumed idle expiry did not close the exact owner") + } +} + +func TestQuiesceStopsIdleTimerUntilResume(t *testing.T) { + sender := suspendSender(func(context.Context, proto.Envelope) error { return nil }) + r := suspensionRouter(t, sender) + closed := make(chan struct{}) + native := &suspendedExecutor{close: func(context.Context) error { close(closed); return nil }} + owner := &executorState{id: "executor", sessionID: "session", environmentID: "env", native: native, cancel: func() {}} r.mu.Lock() - remaining := len(r.executors) + r.idleTimeout = 40 * time.Millisecond + r.executors[owner.sessionID] = owner + r.scheduleExecutorIdleLocked(owner) r.mu.Unlock() - if remaining != 0 || native.closed.Load() != 1 { - t.Fatal("resume retained or reclosed a retired native owner") + request := proto.EnvironmentSuspendPayload{EnvironmentID: "env", SuspendID: "attempt"} + if err := r.Quiesce(t.Context(), request); err != nil { + t.Fatal(err) + } + select { + case <-closed: + t.Fatal("parked owner's idle timer closed native resources") + case <-time.After(2 * r.idleTimeout): + } + if err := r.Resume(request, sender); err != nil { + t.Fatal(err) + } + select { + case <-closed: + case <-time.After(time.Second): + t.Fatal("resume did not restore idle expiry") } } func TestShutdownDestroysQuiescedOwnerAndCannotResume(t *testing.T) { sender := suspendSender(func(context.Context, proto.Envelope) error { return nil }) r := suspensionRouter(t, sender) + native := &suspendedExecutor{} + r.executors["session"] = &executorState{id: "executor", sessionID: "session", environmentID: "env", native: native, cancel: func() {}} request := proto.EnvironmentSuspendPayload{EnvironmentID: "env", SuspendID: "attempt"} if err := r.Quiesce(context.Background(), request); err != nil { t.Fatal(err) @@ -170,6 +216,9 @@ func TestShutdownDestroysQuiescedOwnerAndCannotResume(t *testing.T) { if !errors.Is(r.Resume(request, sender), ErrRouterClosed) { t.Fatal("closed Router resurrected") } + if native.closed.Load() != 1 || len(r.executors) != 0 { + t.Fatal("shutdown did not retire the retained native owner") + } } func TestQuiesceDrainDeadlineCannotReopenAdmission(t *testing.T) { @@ -191,47 +240,17 @@ func TestQuiesceDrainDeadlineCannotReopenAdmission(t *testing.T) { } } -func TestQuiesceWaitsForIdleExecutorClose(t *testing.T) { - entered, release := make(chan struct{}), make(chan struct{}) - native := &suspendedExecutor{close: func(context.Context) error { close(entered); <-release; return nil }} - r := suspensionRouter(t, suspendSender(func(context.Context, proto.Envelope) error { return nil })) - t.Cleanup(func() { - select { - case <-release: - default: - close(release) - } - }) - owner := &executorState{id: "executor", sessionID: "session", environmentID: "env", native: native, cancel: func() {}} - r.executors[owner.sessionID] = owner - quiet := make(chan error, 1) - go func() { - quiet <- r.Quiesce(t.Context(), proto.EnvironmentSuspendPayload{EnvironmentID: "env", SuspendID: "attempt"}) - }() - <-entered - select { - case err := <-quiet: - t.Fatalf("acknowledged before native close: %v", err) - case <-time.After(20 * time.Millisecond): - } - close(release) - if err := <-quiet; err != nil { - t.Fatal(err) - } - if native.closed.Load() != 1 { - t.Fatal("native Close was duplicated") - } -} - -func TestQuiesceCloseFailureRetainsOwnerAndFencesAdmission(t *testing.T) { - failure := errors.New("native history flush failed") +func TestQuiescedShutdownRetainsFailedCloseForRetry(t *testing.T) { + failure := errors.New("native cleanup failed") native := &suspendedExecutor{close: func(context.Context) error { return failure }} r := suspensionRouter(t, suspendSender(func(context.Context, proto.Envelope) error { return nil })) owner := &executorState{id: "executor", sessionID: "session", environmentID: "env", native: native, cancel: func() {}} r.executors[owner.sessionID] = owner - request := proto.EnvironmentSuspendPayload{EnvironmentID: "env", SuspendID: "attempt"} - if err := r.Quiesce(t.Context(), request); !errors.Is(err, failure) || !errors.Is(err, ErrRouterQuiesced) { - t.Fatalf("quiesce lost native close failure: %v", err) + if err := r.Quiesce(t.Context(), proto.EnvironmentSuspendPayload{EnvironmentID: "env", SuspendID: "attempt"}); err != nil || native.closed.Load() != 0 { + t.Fatalf("quiesce retired owner: err=%v closes=%d", err, native.closed.Load()) + } + if err := r.Shutdown(t.Context()); !errors.Is(err, failure) { + t.Fatalf("shutdown lost native close failure: %v", err) } r.mu.Lock() retained := r.executors[owner.sessionID] == owner && owner.invalid && owner.closeErr == failure @@ -239,9 +258,6 @@ func TestQuiesceCloseFailureRetainsOwnerAndFencesAdmission(t *testing.T) { if !retained { t.Fatal("failed close abandoned ownership") } - if err := r.Handle(t.Context(), proto.Envelope{Type: proto.TypeExecutionPrepare, ID: "late"}); !errors.Is(err, ErrRouterQuiesced) { - t.Fatalf("failed close reopened admission: %v", err) - } native.close = nil if err := r.Shutdown(t.Context()); err != nil { t.Fatal(err) @@ -251,7 +267,7 @@ func TestQuiesceCloseFailureRetainsOwnerAndFencesAdmission(t *testing.T) { } } -func TestQuiesceDeadlineJoinsNativeCloseDuringShutdown(t *testing.T) { +func TestQuiesceCancelledDrainRetainsOwnerUntilShutdownSettles(t *testing.T) { entered, release := make(chan struct{}), make(chan struct{}) native := &suspendedExecutor{close: func(context.Context) error { close(entered); <-release; return nil }} r := suspensionRouter(t, suspendSender(func(context.Context, proto.Envelope) error { return nil })) @@ -264,21 +280,22 @@ func TestQuiesceDeadlineJoinsNativeCloseDuringShutdown(t *testing.T) { }) owner := &executorState{id: "executor", sessionID: "session", environmentID: "env", native: native, cancel: func() {}} r.executors[owner.sessionID] = owner + r.shutdownWG.Add(1) ctx, cancel := context.WithCancel(t.Context()) - quiet := make(chan error, 1) - go func() { - quiet <- r.Quiesce(ctx, proto.EnvironmentSuspendPayload{EnvironmentID: "env", SuspendID: "attempt"}) - }() - <-entered cancel() - if err := <-quiet; !errors.Is(err, context.Canceled) { + if err := r.Quiesce(ctx, proto.EnvironmentSuspendPayload{EnvironmentID: "env", SuspendID: "attempt"}); !errors.Is(err, context.Canceled) { t.Fatalf("quiesce=%v", err) } + if native.closed.Load() != 0 { + t.Fatal("cancelled drain closed resident owner") + } if err := r.Handle(t.Context(), proto.Envelope{Type: proto.TypeExecutionPrepare, ID: "late"}); !errors.Is(err, ErrRouterQuiesced) { t.Fatalf("deadline reopened admission: %v", err) } + r.shutdownWG.Done() stopped := make(chan error, 1) go func() { stopped <- r.Shutdown(t.Context()) }() + <-entered select { case err := <-stopped: t.Fatalf("shutdown abandoned in-flight close: %v", err) @@ -308,3 +325,44 @@ func TestQuiesceRejectsAdmissionInProgress(t *testing.T) { t.Fatal("busy rejection changed suspension") } } + +func TestQuiesceFencesAlreadyRunningTimerAcrossResume(t *testing.T) { + sender := suspendSender(func(context.Context, proto.Envelope) error { return nil }) + r := suspensionRouter(t, sender) + native := &suspendedExecutor{} + owner := &executorState{id: "executor", sessionID: "session", environmentID: "env", native: native, cancel: func() {}} + entered, release, done := make(chan struct{}), make(chan struct{}), make(chan struct{}) + t.Cleanup(func() { + select { + case <-release: + default: + close(release) + } + <-done + }) + r.mu.Lock() + r.executors[owner.sessionID] = owner + r.scheduleExecutorIdleLocked(owner) + owner.timer.Stop() + lease := owner.idleLease + owner.timer = time.AfterFunc(0, func() { + close(entered) + <-release + r.expireIdleExecutor(owner, lease) + close(done) + }) + r.mu.Unlock() + <-entered + request := proto.EnvironmentSuspendPayload{EnvironmentID: "env", SuspendID: "attempt"} + if err := r.Quiesce(t.Context(), request); err != nil { + t.Fatal(err) + } + if err := r.Resume(request, sender); err != nil { + t.Fatal(err) + } + close(release) + <-done + if native.closed.Load() != 0 { + t.Fatal("callback already running before suspension retired resumed owner") + } +} diff --git a/contracts/agents-api/harness-onboarding.md b/contracts/agents-api/harness-onboarding.md index 6032efbe2..d473444ef 100644 --- a/contracts/agents-api/harness-onboarding.md +++ b/contracts/agents-api/harness-onboarding.md @@ -102,7 +102,7 @@ The service profile qualifies public combinations and the Runtime advertises the A Session owns one reusable Executor in its connected Runtime; a Turn owns one input execution, its output stream and its cancellation. `agent.ExecutorFactory` prepares the fixed configuration without model input, and `Executor.StartTurn` creates a new `agent.Turn` without replacing healthy native resources. Normal completion settles only the Turn. `Executor.Close` releases native resources on idle expiry, Environment shutdown or confirmed invalidation; it releases neither the Environment allocation nor the workspace. Core keeps no second Executor cache. The same lifecycle applies to hosted, self-hosted and `none` placements. -**Binding.** The Runtime binds its Executor record to the Session, Environment, connection and immutable execution configuration. Resume identity and prior-Turn recovery flags are continuity assertions, not configuration changes. A supplied native identity must match the retained owner, and when existing history is required, recovery never starts a new root. A configuration conflict is an error, not a hot switch. A lost connection retires its owners and handles; old timers, output and cancellation cannot affect their replacements. +**Binding.** The Runtime binds its Executor record to the Session, Environment, connection and immutable execution configuration. Resume identity and prior-Turn recovery flags are continuity assertions, not configuration changes. A supplied native identity must match the retained owner, and when existing history is required, recovery never starts a new root. A configuration conflict is an error, not a hot switch. An ordinary lost connection retires its owners and handles; [planned suspension](../../docs/runtime-protocol.md#preparation-and-execution-order) retains settled idle owners through exact authenticated resume. Old timers, output and cancellation cannot affect replacement owners. **Per-Turn state.** Each Turn gets a fresh wrapper, output channel and receipt state. Steering and function interfaces belong to that Turn. Native callbacks capture the originating Turn before asynchronous work, so a late event is never attributed to whichever Turn is active. Native processes, query or transport connections, fixed capability configuration and native session identity belong to the Executor. Do not reset completed `sync.Once` values or reuse an old Turn object. diff --git a/contracts/agents-api/zh/harness-onboarding.md b/contracts/agents-api/zh/harness-onboarding.md index 5df702061..d422a7116 100644 --- a/contracts/agents-api/zh/harness-onboarding.md +++ b/contracts/agents-api/zh/harness-onboarding.md @@ -1,7 +1,7 @@ --- title: "添加 Harness" source: contracts/agents-api/harness-onboarding.md -source_hash: 779954060b3be129f858e2f57e62b1af56cfe2f925f98362b48b20d32574f681 +source_hash: 2fe919ec4fd2285cc9e54fed58ef66aebe66d132fd8fb5adc272e26b2e673774 --- **Harness** 是一种运行模型和工具循环的原生代理引擎(Codex、Claude Code、MiniMax Code)。**Harness 适配器**将 Runtime 的 Executor 和 Turn 契约转换到该引擎的 SDK 或协议。本文档定义 Runtime–Harness 协议:适配器接口及其生命周期义务、注册、Core 资格认定和验收。[Harness capabilities](harness-capabilities.md) 记录了当前每个 Harness 支持的功能。 @@ -104,7 +104,7 @@ func (s *Session) SubmitFunctionResult(context.Context, proto.FunctionResultPayl Session 在其已连接的 Runtime 中拥有一个可复用的 Executor;Turn 拥有一次输入执行、其输出流和其取消操作。`agent.ExecutorFactory` 在没有模型输入的情况下准备固定配置,而 `Executor.StartTurn` 创建新的 `agent.Turn`,不替换健康的原生资源。正常完成仅结算 Turn。`Executor.Close` 在空闲过期、Environment 关闭或确认失效时释放原生资源;它既不释放 Environment 分配,也不释放工作区。Core 不保留第二套 Executor 缓存。相同的生命周期适用于托管、自托管和 `none` 放置方式。 -**绑定。** Runtime 将其 Executor 记录绑定到 Session、Environment、连接和不可变执行配置。恢复身份和先前 Turn 恢复标志是连续性断言,而不是配置更改。提供的原生身份必须与保留的所有者匹配;当需要现有历史时,恢复绝不能启动新的根。配置冲突属于错误,而不是热切换。连接丢失会让其所有者和句柄退役;旧计时器、输出和取消操作不能影响替代对象。 +**绑定。** Runtime 将其 Executor 记录绑定到 Session、Environment、连接和不可变执行配置。恢复身份和先前 Turn 恢复标志是连续性断言,而不是配置更改。提供的原生身份必须与保留的所有者匹配;当需要现有历史时,恢复绝不能启动新的根。配置冲突属于错误,而不是热切换。普通连接丢失会让其所有者和句柄退役;[计划性挂起](../../../docs/zh/runtime-protocol.md#preparation-and-execution-order) 则保留已结算的空闲所有者,直至精确匹配且已认证的恢复完成。旧计时器、输出和取消操作不能影响替代所有者。 **每 Turn 状态。** 每个 Turn 都会获得全新的包装器、输出通道和回执状态。引导和函数接口均属于该 Turn。原生回调必须在异步工作开始前捕获来源 Turn,因此迟到事件绝不会被归到当前活动的 Turn 上。原生进程、query 或传输连接、固定能力配置和原生 session 身份均属于 Executor。不要重置已完成的 `sync.Once` 值,也不要复用旧 Turn 对象。 diff --git a/docs/runtime-protocol.md b/docs/runtime-protocol.md index e7c818fa8..f892f27a6 100644 --- a/docs/runtime-protocol.md +++ b/docs/runtime-protocol.md @@ -126,7 +126,7 @@ Preparation and start run outside the receive loop and router lock. An admission Idle expiry of an Executor is a Runtime resource policy, separate from Core's active-Turn concurrency. On shutdown the Runtime closes active and idle Executors, keeps any target whose close failed and allows a later serialized retry. An ordinary disconnection closes the failed transport and keeps the exact router until shutdown succeeds; a wait timeout or failed cleanup never authorizes reconnection, and process shutdown keeps waiting rather than discarding owned native resources. Workspace operations keep their binding and settlement rules across Turn boundaries and Executor closure. -Suspension closes admission, drains admitted work and receipts, and confirms closure of every idle Executor before acknowledging `environment_quiesced`. A failed or unconfirmed native close prevents acknowledgement and retains cleanup ownership. The Sandbox Provider remains responsible for the filesystem flush and compute-stop guarantees of its checkpoint implementation; Runtime quiescence alone does not establish those guarantees. +Suspension closes admission and drains admitted work and receipts while retaining settled idle Executors before acknowledging `environment_quiesced`. Active, preparing, invalid or otherwise unsettled owners prevent suspension. Idle expiry is stopped and its callbacks fenced throughout suspension; exact authenticated `environment_resume` restores a fresh idle interval on the same owners. Only this planned reconnect preserves the Router and native owners across sockets; ordinary disconnect, shutdown and cancellation of the connection lifecycle still require confirmed cleanup. A drain failure keeps admission closed until shutdown. Preparation after resume retains the same native identity and immutable-configuration checks, including credential conflicts; it never silently replaces a conflicting resident Executor. The Sandbox Provider remains responsible for the filesystem flush and compute-stop guarantees of its checkpoint implementation; Runtime quiescence alone proves neither those guarantees nor restoration of native RAM or external network connections. ## Active input receipts diff --git a/docs/zh/runtime-protocol.md b/docs/zh/runtime-protocol.md index ab1b54296..279231ac7 100644 --- a/docs/zh/runtime-protocol.md +++ b/docs/zh/runtime-protocol.md @@ -1,7 +1,7 @@ --- title: "Core–Runtime 协议" source: docs/runtime-protocol.md -source_hash: 142a3d2a3dc5b790bcc80dff2364d094857f835e0e0687f91373251a4f4d9b3b +source_hash: 902f53e5471de13ecac0e914e4d23d7650955a4c50ec1aa01517c13aef7946f7 --- 此协议在 Runtime daemon 获取机器凭据后连接 Core 与 daemon,定义 daemon 连接上消息的含义和顺序。wire 类型、限制和验证器仅在 [`internal/agentdaemon/proto`](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/internal/agentdaemon/proto) 中定义一次;Core 的 [gateway](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/services/core/internal/runtimegateway) 与参考 Runtime 的 [dispatcher](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/apps/daemon/internal/dispatch) 都使用它们,因此无需同步第二套 payload schema。签发凭据和打开连接的 HTTP 路由见[机器连接 API](../../contracts/agents-api/zh/machine-api.md)。 @@ -129,7 +129,7 @@ preparation 和 start 在 receive loop 与 router lock 之外运行。admission Executor 空闲到期属于 Runtime 资源策略,与 Core 的活动 Turn 并发限制独立。关闭时 Runtime 关闭活动和空闲 Executor,保留关闭失败的目标,并允许稍后串行重试。普通断连会关闭失败的 transport 并保留原 router,直到 shutdown 成功;等待超时或清理失败不授权重连,进程 shutdown 继续等待,不丢弃自己拥有的原生资源。工作区操作在跨 Turn 和 Executor 关闭后仍保留绑定与结算规则。 -挂起先关闭准入、排空已接纳工作与回执,并确认每个空闲 Executor 已关闭,随后才确认 `environment_quiesced`。原生关闭失败或结果未确认时,不得确认挂起,且保留清理责任。Sandbox Provider 仍负责其 checkpoint 实现的文件系统刷新和计算停止保证;Runtime 静止本身不证明这些保证。 +挂起先关闭准入、排空已接纳工作与回执,同时保留已结算的空闲 Executor,随后才确认 `environment_quiesced`。活动、准备中、失效或其他未结算的所有者会阻止挂起。整个挂起期间停止空闲到期计时并隔离其回调;精确匹配且已认证的 `environment_resume` 为同一批所有者恢复一个新的空闲计时间隔。只有这种计划性重连跨 socket 保留 Router 与原生所有者;普通断连、shutdown 和连接生命周期取消仍要求确认清理完成。排空失败后,准入保持关闭直到 shutdown。恢复后的准备仍检查相同的原生身份与不可变配置,包括凭据冲突;绝不静默替换配置冲突的驻留 Executor。Sandbox Provider 仍负责其 checkpoint 实现的文件系统刷新和计算停止保证;Runtime 静止本身既不证明这些保证,也不证明原生 RAM 或外部网络连接已恢复。 ## 活动输入回执 {#active-input-receipts} diff --git a/internal/agentdaemon/proto/envelope_test.go b/internal/agentdaemon/proto/envelope_test.go index d485d8f7d..cec6e5d6d 100644 --- a/internal/agentdaemon/proto/envelope_test.go +++ b/internal/agentdaemon/proto/envelope_test.go @@ -88,6 +88,7 @@ func TestVersionCompatible(t *testing.T) { ok bool }{ {Version, true}, // exact match + {"0.15.0", false}, // closes idle Executors before suspension {"0.8.99", false}, // patch drift NOT OK {"0.6.99", false}, // minor drift NOT OK {"1.0.0", false}, // major drift NOT OK diff --git a/internal/agentdaemon/proto/suspend.go b/internal/agentdaemon/proto/suspend.go index 0e0ad7896..b9c4dede7 100644 --- a/internal/agentdaemon/proto/suspend.go +++ b/internal/agentdaemon/proto/suspend.go @@ -18,8 +18,9 @@ type EnvironmentSuspendPayload struct { } // An accepted quiesce confirms closed admission, drained work and receipts, -// and successful closure of idle native Executors. It does not confirm a -// filesystem checkpoint or that provider-owned compute has stopped. +// and retention of settled idle native Executors with idle expiry fenced until +// exact authenticated resume. It does not confirm a filesystem checkpoint or +// that provider-owned compute has stopped. type EnvironmentSuspendResultPayload struct { EnvironmentID string `json:"environment_id"` SuspendID string `json:"suspend_id"` diff --git a/internal/agentdaemon/proto/version.go b/internal/agentdaemon/proto/version.go index 114299f22..2dcbf7337 100644 --- a/internal/agentdaemon/proto/version.go +++ b/internal/agentdaemon/proto/version.go @@ -2,7 +2,7 @@ package proto // Version identifies the complete Core–Runtime wire contract. Change it when // removing or changing a payload or its semantics; deploy both endpoints together. -const Version = "0.15.0" +const Version = "0.16.0" // VersionCompatible accepts only this contract. Patch drift, prerelease suffixes // and malformed versions do not select an implicit compatibility path.