From 50b53858a157b6bbd276797e2c84150197279bac Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Fri, 9 Oct 2026 00:33:26 +0000 Subject: [PATCH] Keep workspace reads independent of native readiness --- .../internal/dispatch/prepared_handoff.go | 4 - .../dispatch/prepared_handoff_test.go | 18 ++- .../dispatch/workspace_directory_test.go | 35 +++-- .../internal/dispatch/workspace_export.go | 2 +- .../internal/dispatch/workspace_read.go | 16 +- docs/runtime-protocol.md | 6 +- docs/zh/runtime-protocol.md | 8 +- internal/agentdaemon/proto/version.go | 2 +- .../agentdaemon/proto/workspace_directory.go | 2 +- internal/agentdaemon/proto/workspace_read.go | 3 +- .../execution/environment_directory.go | 14 -- services/core/internal/execution/worker.go | 13 +- .../internal/execution/worker_schedule.go | 13 +- .../environment_directory_active_test.go | 11 +- ...ronment_directory_preparing_public_test.go | 145 ++++++++++++++++++ .../integration/environment_directory_test.go | 20 ++- .../integration/local_artifact_export_test.go | 2 +- .../local_environment_file_write_test.go | 9 +- .../local_environment_worker_test.go | 2 +- 19 files changed, 246 insertions(+), 79 deletions(-) create mode 100644 services/core/tests/integration/environment_directory_preparing_public_test.go diff --git a/apps/daemon/internal/dispatch/prepared_handoff.go b/apps/daemon/internal/dispatch/prepared_handoff.go index e5fdf3856..3b2b28253 100644 --- a/apps/daemon/internal/dispatch/prepared_handoff.go +++ b/apps/daemon/internal/dispatch/prepared_handoff.go @@ -89,10 +89,6 @@ func (r *Router) preparedOperationLocked(state *sessionState) (agent.Turn, func( return state.session, handoff.operations.RUnlock, true } -func (r *Router) runRouteOpenLocked(state *sessionState) bool { - return state != nil && !r.closed && state.ctx.Err() == nil && !state.steeringClosed && state.preparedHandoff.release == nil -} - // claimPreparedReleaseLocked closes admission permanently and returns the // current native attempt. retry starts a new serialized attempt only after the // previous one failed. Router.mu must be held. diff --git a/apps/daemon/internal/dispatch/prepared_handoff_test.go b/apps/daemon/internal/dispatch/prepared_handoff_test.go index 89e40d832..93c18b55a 100644 --- a/apps/daemon/internal/dispatch/prepared_handoff_test.go +++ b/apps/daemon/internal/dispatch/prepared_handoff_test.go @@ -180,15 +180,29 @@ func TestPreparedHandoffDuplicateStartDoesNotReexecuteDuringPublication(t *testi if ack := lastSteeringAck(t, sender.recSender, "run", "publication-steering"); ack.ErrorCode != "not_ready" { t.Fatalf("pre-publication steering = %+v", ack) } - read := proto.WorkspaceReadPayload{RunID: "run", EnvironmentID: preparationEnvironmentID, MaxEntries: 1} + read := map[string]any{"run_id": "run", "environment_id": preparationEnvironmentID, "max_entries": 1} if err := r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceRead, "publication-read", read)); err != nil { t.Fatal(err) } var readResult proto.WorkspaceReadResultPayload frames := sender.snapshot() - if len(frames) == 0 || frames[len(frames)-1].Type != proto.TypeWorkspaceReadResult || frames[len(frames)-1].DecodePayload(&readResult) != nil || readResult.ErrorCode != "resource_unavailable" { + if len(frames) == 0 || frames[len(frames)-1].Type != proto.TypeWorkspaceReadResult || frames[len(frames)-1].DecodePayload(&readResult) != nil || readResult.ErrorCode != "invalid_request" { t.Fatalf("pre-publication workspace read = %+v", frames) } + readPreparation := preparationRequest() + readPreparation.Configuration = proto.PromptRequestPayload{AgentKind: "prepared", LocalEnvironment: &proto.LocalEnvironment{ID: preparationEnvironmentID}, WorkspaceReadOnly: true} + if err := r.Handle(t.Context(), mustEnv(t, proto.TypeExecutionPrepare, "publication-read-prepare", readPreparation)); err != nil { + t.Fatal(err) + } + readReady := waitPreparationStatus(t, sender.recSender, "publication-read-prepare", "ready", "") + if err := r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceRead, "publication-independent-read", proto.WorkspaceReadPayload{Handle: readReady.Handle, EnvironmentID: preparationEnvironmentID, MaxEntries: 1})); err != nil { + t.Fatal(err) + } + if got := waitWorkspaceRead(t, sender.recSender, "publication-independent-read"); got.Outcome != "completed" || !got.CloseAcknowledged || got.Directory == nil { + t.Fatal("read-only preparation depends on started publication", got) + } + _ = r.Handle(t.Context(), mustEnv(t, proto.TypeExecutionRelease, "publication-read-prepare", proto.ExecutionReleasePayload{Handle: readReady.Handle})) + waitPreparationStatus(t, sender.recSender, "publication-read-prepare", "released", "") if session.functions.Load() != 0 || session.steers.Load() != 0 { t.Fatal("private Session accepted work before started publication") } diff --git a/apps/daemon/internal/dispatch/workspace_directory_test.go b/apps/daemon/internal/dispatch/workspace_directory_test.go index 8c1adb1fd..c74f21ecc 100644 --- a/apps/daemon/internal/dispatch/workspace_directory_test.go +++ b/apps/daemon/internal/dispatch/workspace_directory_test.go @@ -9,7 +9,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) -func TestWorkspaceDirectoryRetainsEnvironmentAndTransferredOwner(t *testing.T) { +func TestWorkspaceDirectoryUsesReadOnlyOwnerBesideNativeExecution(t *testing.T) { sender := &recSender{} p := &controlledPreparation{closed: make(chan struct{})} p.start = func(ctx context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { @@ -21,7 +21,13 @@ func TestWorkspaceDirectoryRetainsEnvironmentAndTransferredOwner(t *testing.T) { r := ownedPreparationRouter(t, sender, time.Minute, func(context.Context, proto.PromptRequestPayload) (preparedFixture, error) { return p, nil }, owner) _ = r.Handle(t.Context(), mustEnv(t, proto.TypeExecutionPrepare, "prepare", preparationRequest())) ready := waitPreparationStatus(t, sender, "prepare", "ready", "") - request := proto.WorkspaceReadPayload{Handle: ready.Handle, EnvironmentID: preparationEnvironmentID, MaxEntries: 1} + configuration := preparationRequest() + configuration.Configuration = proto.PromptRequestPayload{AgentKind: "prepared", LocalEnvironment: &proto.LocalEnvironment{ID: preparationEnvironmentID}, WorkspaceReadOnly: true} + if err := r.Handle(t.Context(), mustEnv(t, proto.TypeExecutionPrepare, "read-prepare", configuration)); err != nil { + t.Fatal(err) + } + readReady := waitPreparationStatus(t, sender, "read-prepare", "ready", "") + request := proto.WorkspaceReadPayload{Handle: readReady.Handle, EnvironmentID: preparationEnvironmentID, MaxEntries: 1} for _, phase := range []string{"idle", "active"} { _ = r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceRead, phase, request)) result := waitWorkspaceRead(t, sender, phase) @@ -34,19 +40,22 @@ func TestWorkspaceDirectoryRetainsEnvironmentAndTransferredOwner(t *testing.T) { if got := waitWorkspaceRead(t, sender, phase+"-foreign"); got.ErrorCode != proto.AssignmentConflict { t.Fatal(got) } + native := request + native.Handle = ready.Handle + _ = r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceRead, phase+"-native", native)) + if got := waitWorkspaceRead(t, sender, phase+"-native"); got.ErrorCode != "resource_unavailable" { + t.Fatal("execution preparation admitted a File read", got) + } if phase == "idle" { _ = r.Handle(t.Context(), mustEnv(t, proto.TypeExecutionStart, "prepare", proto.ExecutionStartPayload{Handle: ready.Handle, ExecutorID: ready.ExecutorID, RunID: "run", Input: proto.TextInput("start")})) waitPreparationStatus(t, sender, "prepare", "started", "") - _ = r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceRead, "stale", request)) - if got := waitWorkspaceRead(t, sender, "stale"); got.Outcome != "rejected" { - t.Fatal(got) - } - request.Handle, request.RunID = "", "run" } } - for index, bad := range []proto.WorkspaceReadPayload{ - {RunID: "run", EnvironmentID: preparationEnvironmentID, MaxEntries: proto.WorkspaceDirectoryMaxEntries + 1}, - {Handle: ready.Handle, RunID: "run", EnvironmentID: preparationEnvironmentID, MaxEntries: 2}, + for index, bad := range []any{ + proto.WorkspaceReadPayload{Handle: readReady.Handle, EnvironmentID: preparationEnvironmentID, MaxEntries: proto.WorkspaceDirectoryMaxEntries + 1}, + proto.WorkspaceReadPayload{EnvironmentID: preparationEnvironmentID, MaxEntries: 2}, + map[string]any{"run_id": "run", "environment_id": preparationEnvironmentID, "max_entries": 2}, + map[string]any{"handle": readReady.Handle, "run_id": "run", "environment_id": preparationEnvironmentID, "max_entries": 2}, } { id := fmt.Sprintf("invalid-%d", index) _ = r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceRead, id, bad)) @@ -54,4 +63,10 @@ func TestWorkspaceDirectoryRetainsEnvironmentAndTransferredOwner(t *testing.T) { t.Fatal("malformed directory control reached a resource", got) } } + _ = r.Handle(t.Context(), mustEnv(t, proto.TypeExecutionRelease, "read-prepare", proto.ExecutionReleasePayload{Handle: readReady.Handle})) + waitPreparationStatus(t, sender, "read-prepare", "released", "") + _ = r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceRead, "released-read", request)) + if got := waitWorkspaceRead(t, sender, "released-read"); got.ErrorCode != "resource_unavailable" { + t.Fatal("released read owner remained usable", got) + } } diff --git a/apps/daemon/internal/dispatch/workspace_export.go b/apps/daemon/internal/dispatch/workspace_export.go index e7271f30e..d26ee4003 100644 --- a/apps/daemon/internal/dispatch/workspace_export.go +++ b/apps/daemon/internal/dispatch/workspace_export.go @@ -49,7 +49,7 @@ func (r *Router) handleWorkspaceExport(ctx context.Context, env proto.Envelope) } environment, code := r.workspaceResourceLocked(env.Assignment, proto.WorkspaceReadPayload{Handle: request.Handle, EnvironmentID: request.EnvironmentID}) p := r.preparations[request.Handle] - if code == "" && (u != nil || r.workspaceWrites[env.Assignment.SessionID] != nil || p.executor != nil) { + if code == "" && (u != nil || r.workspaceWrites[env.Assignment.SessionID] != nil) { code = "resource_unavailable" } if code != "" { diff --git a/apps/daemon/internal/dispatch/workspace_read.go b/apps/daemon/internal/dispatch/workspace_read.go index ad0b3f357..1cbe20597 100644 --- a/apps/daemon/internal/dispatch/workspace_read.go +++ b/apps/daemon/internal/dispatch/workspace_read.go @@ -61,7 +61,7 @@ func (r *Router) handleWorkspaceRead(ctx context.Context, env proto.Envelope) er } // workspaceResourceLocked returns the Environment owner of ref's Session once -// ref admits the read's ready preparation handle or open Run. +// ref admits the read-only preparation's ready handle. func (r *Router) workspaceResourceLocked(ref proto.AssignmentRef, request proto.WorkspaceReadPayload) (Environment, string) { if code := r.admitLocked(ref, ref.SessionID, request.EnvironmentID); code != "" { return nil, code @@ -70,17 +70,9 @@ func (r *Router) workspaceResourceLocked(ref proto.AssignmentRef, request proto. if environment == nil { return nil, "read_unsupported" } - if request.Handle != "" { - p := r.preparations[request.Handle] - if p == nil || p.request.Assignment != ref || p.environmentID != request.EnvironmentID || p.status.State != "ready" || - !p.owns || p.busy || p.ctx.Err() != nil || !time.Now().Before(p.deadline) { - return nil, "resource_unavailable" - } - return environment, "" - } - s := r.sessions[request.RunID] - if s == nil || s.assignment != ref || s.environmentID != request.EnvironmentID || s.session == nil || - !r.runRouteOpenLocked(s) { + p := r.preparations[request.Handle] + if p == nil || p.request.Assignment != ref || p.environmentID != request.EnvironmentID || p.status.State != "ready" || + p.executor != nil || !p.owns || p.busy || p.ctx.Err() != nil || !time.Now().Before(p.deadline) { return nil, "resource_unavailable" } return environment, "" diff --git a/docs/runtime-protocol.md b/docs/runtime-protocol.md index 1c81ca2b5..d0d40270d 100644 --- a/docs/runtime-protocol.md +++ b/docs/runtime-protocol.md @@ -194,13 +194,13 @@ Core stores accepted values in the Turn outcome as `engine_error_code` and `engi ## Workspace operations -A workspace read that needs no running Turn uses the read-only preparation profile: `execution_prepare` with `workspace_read_only`, which requires the `local_environment` capability. It accepts only the bound Environment and resource identity; execution options, model and MCP credentials, native Session continuation and model or tool input are excluded, and the owner rejects `execution_start`. The Session's Environment owner serves it without starting a Harness process. The profile publishes `released` after the Runtime drops the preparation's ownership; a stale status snapshot never publishes success. A release request, HTTP disconnect or remote socket closure alone does not confirm the release. +A workspace read uses the read-only preparation profile independently of native preparation or a running Turn: `execution_prepare` with `workspace_read_only`, which requires the `local_environment` capability. It accepts only the bound Environment and resource identity; execution options, model and MCP credentials, native Session continuation and model or tool input are excluded, and the owner rejects `execution_start`. The Session's Environment owner serves it without starting a Harness process. The profile publishes `released` after the Runtime drops the preparation's ownership; a stale status snapshot never publishes success. A release request, HTTP disconnect or remote socket closure alone does not confirm the release. -`workspace_read` lists one workspace-relative directory (an empty path selects the root) of an existing preparation handle, or the Run it was transferred to, on the same authenticated device connection, with the exact frozen Environment identity; callers cannot supply sockets, credentials or workspace roots. A result carries at most `max_entries` (1 to 1024) single-component UTF-8 names of at most 255 bytes each, the entry kind, regular-file sizes and explicit truncation, and is returned only after directory access and handle cleanup settle. There is no snapshot, recursion or pagination at this layer. +`workspace_read` lists one workspace-relative directory (an empty path selects the root) through a ready read-only preparation handle on the same authenticated device connection, with the exact frozen Environment identity; execution preparation handles are rejected, and callers cannot supply Run identities, sockets, credentials or workspace roots. A result carries at most `max_entries` (1 to 1024) single-component UTF-8 names of at most 255 bytes each, the entry kind, regular-file sizes and explicit truncation, and is returned only after directory access and handle cleanup settle. There is no snapshot, recursion or pagination at this layer. Request payloads are bounded at 8 KiB and correlation IDs at 128 bytes before admission; oversized IDs are not echoed, and oversized trace metadata is omitted from replies. These are private transport limits, not public Files parameters. A safe native rejection carries no directory; an interrupted or ambiguous read stays unknown and stops further reads on that owner. A dispatched read keeps its bounded waiter across observer cancellation and resource transfer or release, and resource closure stops new admission. The gateway bounds subscriptions and never retries or replays a read on reconnect; a duplicate pending operation ID cannot start another read. -Core runs an idle directory read on the Worker's Session scheduling reservation and targets the exact Run during active execution. It keeps the reservation through the bounded read and release, returns data only after a confirmed close (an incomplete read or uncertain cleanup returns unavailable without data), releases the reservation before delivering the result, and revokes the scoped read credential on completion or failure. The Runtime keeps uncertain cleanup ownership and capacity. The [Environment Files contract](../contracts/agents-api/environment-files.md) owns public authorization, paths and pagination. +Core uses an independent read-only preparation for every directory read and reserves the Worker's Session scheduling slot when idle. It keeps that reservation through the bounded read and release, returns data only after a confirmed close (an incomplete read or uncertain cleanup returns unavailable without data), releases the reservation before delivering the result, and revokes the scoped read credential on completion or failure. The Runtime keeps uncertain cleanup ownership and capacity. The [Environment Files contract](../contracts/agents-api/environment-files.md) owns public authorization, paths and pagination. `workspace_write` transfers a complete bounded body in acknowledged 64 KiB frames before the native writer runs, verifies the declared digest and runs no model. The private transfer bound is 50 MiB, separate from the public 5 MiB decoded inline bound that the API checks before any Runtime work. The Runtime excludes the Session's execution while it receives or applies a write; a malformed, incomplete or expired transfer never reaches the installer. An exact commit or rejection receipt settles the write. A missing or ambiguous receipt leaves the uncertainty with the [Environment owner](#session-assignments): observer cancellation and local process exit cannot prove that nothing changed. Before public admission Core durably reserves the write under the Session lock and blocks successor mutations across restarts until exact settlement; the request is never replayed. The agent host's owner lists directories, creates files and exports outputs over [File access](./file-access-protocol.md) on its Link attachment. diff --git a/docs/zh/runtime-protocol.md b/docs/zh/runtime-protocol.md index 7e21874a8..e3d5965e5 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: ad72e725702ce94ae55ccae0f26f0df9ba43bf5fbd960998bc764bc808aac882 +source_hash: 5003cac2a7dd5d1926c8004a400a12cb947070581351dc45b9f2f7b7323dd726 --- 此协议在 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)。 @@ -196,13 +196,13 @@ Core 在 Turn outcome 中将接受的值保存为 `engine_error_code` 和 `engin ## 工作区操作 {#workspace-operations} -无需运行 Turn 的工作区读取使用只读 preparation profile:带 `workspace_read_only` 的 `execution_prepare`,要求 `local_environment` 能力。仅接受绑定的 Environment 和 resource 身份;不包含 execution option、model 与 MCP 凭据、原生 Session continuation、model 或 tool 输入,owner 拒绝 `execution_start`。Session 的 Environment owner 提供读取,不启动 Harness 进程。profile 在 Runtime 放弃该 preparation 的所有权后发布 `released`;旧 status snapshot 不发布成功。release 请求、HTTP 断连或远端 socket 关闭本身都不确认释放。 +工作区读取使用独立于原生 preparation 或运行中 Turn 的只读 preparation profile:带 `workspace_read_only` 的 `execution_prepare`,要求 `local_environment` 能力。仅接受绑定的 Environment 和 resource 身份;不包含 execution option、model 与 MCP 凭据、原生 Session continuation、model 或 tool 输入,owner 拒绝 `execution_start`。Session 的 Environment owner 提供读取,不启动 Harness 进程。profile 在 Runtime 放弃该 preparation 的所有权后发布 `released`;旧 status snapshot 不发布成功。release 请求、HTTP 断连或远端 socket 关闭本身都不确认释放。 -`workspace_read` 在同一已认证设备连接上,针对现有 preparation handle 或它已转移给的 Run,使用精确冻结的 Environment 身份,列出一个 workspace 相对目录(空路径选择根目录);调用方不能提供 socket、凭据或 workspace root。结果最多携带 `max_entries`(1 到 1024)个单路径组件 UTF-8 名称,每个最多 255 字节,并包含 entry kind、普通文件大小和明确截断信息;仅在目录访问与 handle 清理结算后返回。此层没有快照、递归或分页。 +`workspace_read` 在同一已认证设备连接上,通过 ready 状态的只读 preparation handle,使用精确冻结的 Environment 身份,列出一个 workspace 相对目录(空路径选择根目录);execution preparation handle 会被拒绝,调用方不能提供 Run 身份、socket、凭据或 workspace root。结果最多携带 `max_entries`(1 到 1024)个单路径组件 UTF-8 名称,每个最多 255 字节,并包含 entry kind、普通文件大小和明确截断信息;仅在目录访问与 handle 清理结算后返回。此层没有快照、递归或分页。 准入前,请求 payload 限制为 8 KiB,correlation ID 限制为 128 字节;过长 ID 不回显,过大的 trace metadata 不进入回复。这些是私有 transport 限制,不是公开 Files 参数。安全的原生拒绝不携带目录;中断或有歧义的读取保持未知,并停止该 owner 的后续读取。已 dispatch 的读取在 observer 取消以及资源转移或释放后仍保留有限时等待方,资源关闭阻止新准入。gateway 限制订阅,不在重连后重试或重放读取;重复的 pending operation ID 不能启动另一次读取。 -Core 在 Worker 的 Session 调度预约上运行空闲目录读取,活动执行时针对精确 Run。它在有限时 read 与 release 期间保留预约,仅在确认 close 后返回数据(不完整读取或不确定清理返回 unavailable,不包含数据),在交付结果前释放预约,并在完成或失败后撤销限定作用域的读取凭据。Runtime 保留不确定清理的所有权和容量。[Environment Files 契约](../../contracts/agents-api/zh/environment-files.md)负责公开授权、路径和分页。 +Core 对每次目录读取都使用独立的只读 preparation,并在 Session 空闲时预约 Worker 的 Session 调度槽位。它在有限时 read 与 release 期间保留该预约,仅在确认 close 后返回数据(不完整读取或不确定清理返回 unavailable,不包含数据),在交付结果前释放预约,并在完成或失败后撤销限定作用域的读取凭据。Runtime 保留不确定清理的所有权和容量。[Environment Files 契约](../../contracts/agents-api/zh/environment-files.md)负责公开授权、路径和分页。 `workspace_write` 在原生 writer 运行前,通过已确认的 64 KiB frame 传输完整且有界的 body,验证声明的 digest,不运行模型。私有 transfer 限制为 50 MiB,与公开 API 在任何 Runtime 工作前检查的 5 MiB decoded inline 限制独立。Runtime 在接收或应用写入时排除该 Session 的执行;格式错误、不完整或到期的 transfer 不会到达 installer。精确的 commit 或拒绝回执结算该写入。缺失或有歧义的回执把不确定性留给 [Environment owner](#session-assignments):observer 取消和本地进程退出不能证明没有改变任何内容。公开准入前,Core 在 Session lock 下持久预约写入,跨重启阻止后继 mutation,直到精确结算;请求不重放。agent host 的 owner 在其 Link attachment 上通过 [File access](./file-access-protocol.md) 列举目录、创建文件和导出输出。 diff --git a/internal/agentdaemon/proto/version.go b/internal/agentdaemon/proto/version.go index 524883ee6..7f6f11783 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.13.0" +const Version = "0.14.0" // VersionCompatible accepts only this contract. Patch drift, prerelease suffixes // and malformed versions do not select an implicit compatibility path. diff --git a/internal/agentdaemon/proto/workspace_directory.go b/internal/agentdaemon/proto/workspace_directory.go index d11e3bf3c..1334bf1f7 100644 --- a/internal/agentdaemon/proto/workspace_directory.go +++ b/internal/agentdaemon/proto/workspace_directory.go @@ -36,7 +36,7 @@ func (result *WorkspaceDirectoryResult) UnmarshalJSON(data []byte) error { } func ValidWorkspaceReadRequest(request WorkspaceReadPayload) bool { - return (request.Handle == "") != (request.RunID == "") && request.EnvironmentID != "" && + return request.Handle != "" && request.EnvironmentID != "" && request.MaxEntries >= 1 && request.MaxEntries <= WorkspaceDirectoryMaxEntries } diff --git a/internal/agentdaemon/proto/workspace_read.go b/internal/agentdaemon/proto/workspace_read.go index 63e640bcd..31f3ee20b 100644 --- a/internal/agentdaemon/proto/workspace_read.go +++ b/internal/agentdaemon/proto/workspace_read.go @@ -10,10 +10,9 @@ const ( WorkspaceReadNotDirectory = "not_directory" ) -// WorkspaceReadPayload lists one directory of an existing resource on the current daemon connection. +// WorkspaceReadPayload lists one directory through a read-only preparation on the current daemon connection. type WorkspaceReadPayload struct { Handle string `json:"handle,omitempty"` - RunID string `json:"run_id,omitempty"` EnvironmentID string `json:"environment_id"` Path string `json:"path"` MaxEntries int `json:"max_entries,omitempty"` diff --git a/services/core/internal/execution/environment_directory.go b/services/core/internal/execution/environment_directory.go index d8afd719b..7d945fa99 100644 --- a/services/core/internal/execution/environment_directory.go +++ b/services/core/internal/execution/environment_directory.go @@ -87,21 +87,12 @@ func (w *Worker) runDirectoryRead(owner context.Context, request directoryReadRe result.err = err return } - run := "" - if session.LastTurn != nil && (session.LastTurn.Status == sessions.TurnInProgress || session.LastTurn.Status == sessions.TurnWaiting) { - run = session.LastTurn.ID - } - if (run == "") != reserved { - return - } if reserved { ready, err := w.bindSessionDevice(check, session, func(id string) bool { return w.directoryDeviceReady(check, id, session.Engine) }) if err != nil || !ready { return } } - // Capture retains the public Turn after its native Run has been released. - prepare := reserved || session.LastTurn != nil && session.LastTurn.ArtifactCaptureStarted bound, err := w.dispatcher.SessionsReader.GetSessionDevice(check, session.TenantID, session.ID) if err != nil || bound.SessionEnvironmentID != environment.ID || !w.directoryDeviceReady(check, bound.ID, session.Engine) { return @@ -114,11 +105,6 @@ func (w *Worker) runDirectoryRead(owner context.Context, request directoryReadRe return } read := proto.WorkspaceReadPayload{EnvironmentID: environment.ID, Path: request.path, MaxEntries: proto.WorkspaceDirectoryMaxEntries} - if !prepare { - read.RunID = run - result = readEnvironmentDirectory(owner, peer, bound.Assignment, read) - return - } result = w.dispatcher.readPreparedDirectory(owner, peer, session, environment, bound, read) return } diff --git a/services/core/internal/execution/worker.go b/services/core/internal/execution/worker.go index 5b0cbfe37..523d5cf63 100644 --- a/services/core/internal/execution/worker.go +++ b/services/core/internal/execution/worker.go @@ -158,7 +158,7 @@ func (w *Worker) Run(ctx context.Context) (runErr error) { w.runtimes.drain() w.observeWorkerClosed(closeLease(ctx, w.lease)) }() - active := make(map[string]bool) + active := make(map[string]workerReservation) w.observeSlots(len(active)) type completion struct { id string @@ -208,11 +208,11 @@ func (w *Worker) Run(ctx context.Context) (runErr error) { case err := <-lifecycleDone: return err case request := <-w.fileWrites: - if request.ctx.Err() != nil || active[request.environment.SessionID] || len(active) == w.executionConcurrency() { + if request.ctx.Err() != nil || active[request.environment.SessionID] != 0 || len(active) == w.executionConcurrency() { request.result <- fileWriteResult{err: ErrExecutionUnavailable} continue } - active[request.environment.SessionID] = true + active[request.environment.SessionID] = workspaceReservation w.observeSlots(len(active)) running.Add(1) go func() { @@ -229,13 +229,14 @@ func (w *Worker) Run(ctx context.Context) (runErr error) { } rescanOnCompletion = false case request := <-w.directoryReads: - if request.ctx.Err() != nil || reads == w.executionConcurrency() || (!active[request.environment.SessionID] && len(active) == w.executionConcurrency()) { + reservation := active[request.environment.SessionID] + if request.ctx.Err() != nil || reservation == workspaceReservation || reads == w.executionConcurrency() || (reservation == 0 && len(active) == w.executionConcurrency()) { request.reply(directoryReadResult{err: ErrExecutionUnavailable}) continue } - reserved := !active[request.environment.SessionID] + reserved := reservation == 0 if reserved { - active[request.environment.SessionID] = true + active[request.environment.SessionID] = workspaceReservation w.observeSlots(len(active)) } reads++ diff --git a/services/core/internal/execution/worker_schedule.go b/services/core/internal/execution/worker_schedule.go index bf3a821cc..c995f51c5 100644 --- a/services/core/internal/execution/worker_schedule.go +++ b/services/core/internal/execution/worker_schedule.go @@ -20,7 +20,14 @@ type scheduledWork struct { reservationID string } -func (s *workerSchedule) selectWork(ctx context.Context, w *Worker, devices []string, active map[string]bool) ([]scheduledWork, error) { +type workerReservation uint8 + +const ( + executionReservation workerReservation = iota + 1 + workspaceReservation +) + +func (s *workerSchedule) selectWork(ctx context.Context, w *Worker, devices []string, active map[string]workerReservation) ([]scheduledWork, error) { turns, err := w.dispatcher.SessionsReader.ListExecutionWork(ctx, s.turnCursor, []string{sessions.TurnQueued}, devices) if err != nil { return nil, err @@ -64,7 +71,7 @@ func (s *workerSchedule) selectWork(ctx context.Context, w *Worker, devices []st s.turnCursor = item.TurnID s.environmentFirst = true } - if active[item.SessionID] { + if active[item.SessionID] != 0 { continue } var ready bool @@ -82,7 +89,7 @@ func (s *workerSchedule) selectWork(ctx context.Context, w *Worker, devices []st if !ready { continue } - active[item.SessionID] = true + active[item.SessionID] = executionReservation selected = append(selected, item) } return selected, nil diff --git a/services/core/tests/integration/environment_directory_active_test.go b/services/core/tests/integration/environment_directory_active_test.go index b27f6a15b..8bb8cbc9d 100644 --- a/services/core/tests/integration/environment_directory_active_test.go +++ b/services/core/tests/integration/environment_directory_active_test.go @@ -7,7 +7,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" ) -func TestEnvironmentDirectoryActiveRunUsesExistingOwner(t *testing.T) { +func TestEnvironmentDirectoryActiveRunUsesReadOnlyPreparation(t *testing.T) { h, w, environment := directoryWorker(t) awaitFixtureCapabilities(t, h, workerEnvironmentCapabilities()) pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "execute", []sessions.Input{messageInput("work")}) @@ -24,13 +24,8 @@ func TestEnvironmentDirectoryActiveRunUsesExistingOwner(t *testing.T) { } h.write(prepare.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 3, State: "started", RunID: start.RunID}) result := startDirectoryRead(t.Context(), w, environment) - read := h.read(proto.TypeWorkspaceRead) - var request proto.WorkspaceReadPayload - if read.DecodePayload(&request) != nil || request.RunID != start.RunID || request.Handle != "" || request.EnvironmentID != environment.ID { - t.Fatal("active read selected another execution owner") - } - size := int64(3) - h.write(read.ID, proto.TypeWorkspaceReadResult, proto.WorkspaceReadResultPayload{Outcome: "completed", CloseAcknowledged: true, Directory: &proto.WorkspaceDirectoryResult{Entries: []proto.WorkspaceDirectoryEntry{{Name: "active.txt", Kind: "file", SizeBytes: &size}}}}) + readPrepare, read := prepareDirectoryRead(t, h, environment) + completeDirectoryRead(t, h, readPrepare, read, false, false) if got := awaitDirectoryResult(t, result); got.err != nil || len(got.value.Entries) != 1 { t.Fatal("active read", got.err) } diff --git a/services/core/tests/integration/environment_directory_preparing_public_test.go b/services/core/tests/integration/environment_directory_preparing_public_test.go new file mode 100644 index 000000000..aa2b62ed3 --- /dev/null +++ b/services/core/tests/integration/environment_directory_preparing_public_test.go @@ -0,0 +1,145 @@ +package integration + +import ( + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" + "github.com/google/uuid" +) + +// File reads have their own readiness and must not depend on native readiness. +// Withholding ready/started at the actual gateway fixes each lifecycle boundary; +// the timers below only bound a broken test, not the interleaving. +func TestEnvironmentDirectoryPublicIndependentOfNativeReadiness(t *testing.T) { + for _, phase := range []string{"native_preparing", "started_not_published"} { + t.Run(phase, func(t *testing.T) { + h, w, environment := localWorkerForSession(t, true, `{"agent":{"model":"test-model"},"environment":{"type":"self_hosted","workspace_directory":"/workspace"}}`) + awaitFixtureCapabilities(t, h, workerEnvironmentCapabilities()) + token := uuid.NewString() + auth := newTestAuthenticator(t, []testAPIKey{{OrganizationID: "test-org", ProjectID: h.tenant, SubjectKind: "service_account", SubjectID: "test-runner", TokenSHA256: runtimedevice.HashCredential(token), TenantID: h.tenant}}) + handler, err := publicHandler(t, h.s, auth, "codex", workerExecution(t, w), acceptUnavailable(t)) + if err != nil { + t.Fatal(err) + } + server := httptest.NewServer(handler) + defer server.Close() + defer server.CloseClientConnections() + type response struct { + status int + body, requestID string + err error + } + send := func(method, path, body string) <-chan response { + done := make(chan response, 1) + go func() { + r, err := http.NewRequestWithContext(t.Context(), method, server.URL+path, strings.NewReader(body)) + if err != nil { + done <- response{err: err} + return + } + r.Header.Set("Authorization", "Bearer "+token) + r.Header.Set("OpenAI-Beta", "agents=v1") + r.Header.Set("Content-Type", "application/json") + resp, err := server.Client().Do(r) + if err != nil { + done <- response{err: err} + return + } + defer resp.Body.Close() + data, err := io.ReadAll(resp.Body) + done <- response{resp.StatusCode, string(data), resp.Header.Get("X-Request-ID"), err} + }() + return done + } + await := func(done <-chan response) response { + t.Helper() + select { + case got := <-done: + if got.err != nil { + t.Fatal(got.err) + } + return got + case <-time.After(5 * time.Second): + t.Fatal("public request did not finish at the held lifecycle boundary") + return response{} + } + } + readFiles := func() response { + t.Helper() + files := send(http.MethodGet, "/v1/agents/environments/"+environment.ID+"/files?path=/workspace&limit=100&order=asc", "") + frame := h.read(proto.TypeExecutionPrepare) + var preparation proto.ExecutionPreparePayload + if frame.DecodePayload(&preparation) != nil || !proto.ValidWorkspaceReadPreparation(preparation.Configuration) || preparation.SessionID != h.session.ID || preparation.Configuration.LocalEnvironment.ID != environment.ID { + t.Fatal("file read did not prepare its own exact Environment") + } + handle := acknowledgePreparation(h, frame.ID) + h.write(frame.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 2, State: "ready"}) + read := h.read(proto.TypeWorkspaceRead) + var request proto.WorkspaceReadPayload + if read.DecodePayload(&request) != nil || request.Handle != handle || request.EnvironmentID != environment.ID || request.Path != "" { + t.Fatal("file read lost its independent owner", request) + } + completeDirectoryRead(t, h, frame.ID, read.ID, false, false) + return await(files) + } + admitted := send(http.MethodPost, "/v1/agents/sessions/"+h.session.ID+"/events", `{"events":[{"type":"agent.session.input.message","input":[{"role":"user","content":[{"type":"input_text","text":"work"}]}]}]}`) + prepare := h.read(proto.TypeExecutionPrepare) + handle := acknowledgePreparation(h, prepare.ID) + var got response + if phase == "native_preparing" { + session, err := sessionAdapter(h.s).GetSession(t.Context(), h.tenant, h.session.ID) + if err != nil || session.LastTurn != nil { + t.Fatal("preparing unexpectedly promoted a Turn", err) + } + got = readFiles() + t.Log("preparation held before ready; LastTurn=nil; no native Start sent") + } + h.write(prepare.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 2, State: "ready"}) + frame := h.read(proto.TypeExecutionStart) + var start proto.ExecutionStartPayload + if frame.DecodePayload(&start) != nil || start.RunID == "" { + t.Fatal("execution did not start") + } + if accepted := await(admitted); accepted.status != http.StatusAccepted { + t.Fatalf("message admission: %+v", accepted) + } + if phase == "started_not_published" { + session, err := sessionAdapter(h.s).GetSession(t.Context(), h.tenant, h.session.ID) + if err != nil || session.LastTurn == nil || session.LastTurn.ID != start.RunID || session.LastTurn.Status != sessions.TurnInProgress { + t.Fatal("message admission did not commit the in-progress Turn", err) + } + got = readFiles() + t.Log("events.message=202; Turn=in_progress; started publication held; File read settled through its independent owner") + } + // Settle the execution after the independent read has released ownership. + h.write(prepare.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 3, State: "started", RunID: start.RunID}) + h.write(start.RunID, proto.TypeDone, proto.DonePayload{}) + completeEmptyArtifactExport(t, h) + awaitDaemonRemoteCondition(t, t.Context(), 5*time.Second, "completed Turn after barrier release", func() bool { + turn, err := sessionAdapter(h.s).GetTurn(t.Context(), h.tenant, h.session.ID, start.RunID) + if err != nil { + t.Fatal(err) + } + return turn.Status == sessions.TurnCompleted + }) + assertPreparationReleased(t, h, prepare.ID, handle) + if got.status != http.StatusOK { + t.Errorf("File readiness depends on native readiness: status=%d request_id=%s body=%s", got.status, got.requestID, got.body) + } else { + var page v1.EnvironmentFileList + if err := json.Unmarshal([]byte(got.body), &page); err != nil || page.HasMore || len(page.Data) != 1 || page.Data[0].Path != "/workspace/report.txt" || page.Data[0].SizeBytes != 9 || page.Data[0].EnvironmentID != environment.ID { + t.Errorf("independent File read lost the exact page: %s (decode=%v)", got.body, err) + } + } + }) + } +} diff --git a/services/core/tests/integration/environment_directory_test.go b/services/core/tests/integration/environment_directory_test.go index cb4f64f70..d08a700c5 100644 --- a/services/core/tests/integration/environment_directory_test.go +++ b/services/core/tests/integration/environment_directory_test.go @@ -80,7 +80,7 @@ func prepareDirectoryRead(t *testing.T, h *dispatchHarness, environment sessions h.write(frame.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 2, State: "ready"}) read := h.read(proto.TypeWorkspaceRead) var input proto.WorkspaceReadPayload - if read.DecodePayload(&input) != nil || input.Handle != handle || input.RunID != "" || input.EnvironmentID != environment.ID || input.Path != "reports" { + if read.DecodePayload(&input) != nil || input.Handle != handle || input.EnvironmentID != environment.ID || input.Path != "reports" { t.Fatal("directory request changed binding or path") } return frame.ID, read.ID @@ -251,12 +251,24 @@ func TestEnvironmentDirectoryObserverCancellationRetainsReadOwner(t *testing.T) ctx, cancel := context.WithCancel(t.Context()) result := startDirectoryRead(ctx, w, environment) request, read := prepareDirectoryRead(t, h, environment) + if got := awaitDirectoryResult(t, startDirectoryRead(t.Context(), w, environment)); !errors.Is(got.err, execution.ErrExecutionUnavailable) { + t.Fatal("idle read reservation admitted another reader", got.err) + } cancel() if got := awaitDirectoryResult(t, result); !errors.Is(got.err, execution.ErrExecutionUnavailable) { t.Fatal("cancelled observer result", got.err) } - if _, err := w.ReadEnvironmentDirectory(t.Context(), environment, "reports"); !errors.Is(err, execution.ErrExecutionUnavailable) { - t.Fatal("cancelled observer freed Session owner") + if got := awaitDirectoryResult(t, startDirectoryRead(t.Context(), w, environment)); !errors.Is(got.err, execution.ErrExecutionUnavailable) { + t.Fatal("cancelled observer freed Session owner", got.err) } - completeDirectoryRead(t, h, request, read, false, false) + h.write(read, proto.TypeWorkspaceReadResult, proto.WorkspaceReadResultPayload{Outcome: "completed", CloseAcknowledged: true, Directory: &proto.WorkspaceDirectoryResult{Entries: []proto.WorkspaceDirectoryEntry{}}}) + // Read the very next frame instead of filtering by type: neither rejected + // reader may have sent another preparation while the first owner was held. + _ = h.conn.SetReadDeadline(time.Now().Add(5 * time.Second)) + var frame proto.Envelope + var release proto.ExecutionReleasePayload + if h.conn.ReadJSON(&frame) != nil || frame.Type != proto.TypeExecutionRelease || frame.ID != request || frame.DecodePayload(&release) != nil || release.Handle == "" { + t.Fatal("rejected read dispatched work before releasing the original owner", frame.Type) + } + h.write(request, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: release.Handle, Revision: 3, State: "released"}) } diff --git a/services/core/tests/integration/local_artifact_export_test.go b/services/core/tests/integration/local_artifact_export_test.go index 056591d51..e5c9eec65 100644 --- a/services/core/tests/integration/local_artifact_export_test.go +++ b/services/core/tests/integration/local_artifact_export_test.go @@ -70,7 +70,7 @@ func completeCaptureDirectoryRead(t *testing.T, h *dispatchHarness, worker *exec h.write(frame.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 2, State: "ready"}) read := h.read(proto.TypeWorkspaceRead) var request proto.WorkspaceReadPayload - if read.DecodePayload(&request) != nil || request.Handle != handle || request.RunID != "" || request.EnvironmentID != environment.ID { + if read.DecodePayload(&request) != nil || request.Handle != handle || request.EnvironmentID != environment.ID { t.Fatal("directory read during capture used a finished native Run") } completeDirectoryRead(t, h, frame.ID, read.ID, false, false) diff --git a/services/core/tests/integration/local_environment_file_write_test.go b/services/core/tests/integration/local_environment_file_write_test.go index 730ac4694..45ffc4255 100644 --- a/services/core/tests/integration/local_environment_file_write_test.go +++ b/services/core/tests/integration/local_environment_file_write_test.go @@ -66,13 +66,18 @@ func TestLocalEnvironmentFileWriteOwnsMutationBeforeDispatch(t *testing.T) { if _, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "concurrent", []sessions.Input{messageInput("work")}); !errors.Is(err, sessions.ErrTurnConflict) { t.Fatal("upload admitted concurrent execution", err) } + if got := awaitDirectoryResult(t, startDirectoryRead(t.Context(), w, environment)); !errors.Is(got.err, execution.ErrExecutionUnavailable) { + t.Fatal("write reservation admitted a directory reader", got.err) + } cancel() if got := awaitLocalWrite(t, done); !errors.Is(got.err, execution.ErrExecutionUnavailable) { t.Fatal("detached observer", got.err) } h.write(begin.ID, proto.TypeWorkspaceWriteResult, proto.WorkspaceWriteResultPayload{Outcome: "ready"}) - chunk := h.read(proto.TypeWorkspaceWrite) - if chunk.ID != begin.ID || chunk.DecodePayload(&request) != nil || request.Step != "chunk" || string(request.Data) != "abc" { + // The rejected read cannot send a preparation ahead of the write's chunk. + _ = h.conn.SetReadDeadline(time.Now().Add(5 * time.Second)) + var chunk proto.Envelope + if h.conn.ReadJSON(&chunk) != nil || chunk.Type != proto.TypeWorkspaceWrite || chunk.ID != begin.ID || chunk.DecodePayload(&request) != nil || request.Step != "chunk" || string(request.Data) != "abc" { t.Fatal("body changed") } h.write(begin.ID, proto.TypeWorkspaceWriteResult, proto.WorkspaceWriteResultPayload{Outcome: "received", Offset: 3}) diff --git a/services/core/tests/integration/local_environment_worker_test.go b/services/core/tests/integration/local_environment_worker_test.go index 8e42f1522..9b2361d11 100644 --- a/services/core/tests/integration/local_environment_worker_test.go +++ b/services/core/tests/integration/local_environment_worker_test.go @@ -70,7 +70,7 @@ func TestLocalEnvironmentWorkerDirectoryUsesExactAuthorityWithoutModel(t *testin h.write(frame.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 2, State: "ready"}) read := h.read(proto.TypeWorkspaceRead) var input proto.WorkspaceReadPayload - if read.DecodePayload(&input) != nil || input.EnvironmentID != environment.ID || input.Handle != handle || input.RunID != "" { + if read.DecodePayload(&input) != nil || input.EnvironmentID != environment.ID || input.Handle != handle { t.Fatal("local directory owner changed") } completeDirectoryRead(t, h, frame.ID, read.ID, false, false)