diff --git a/apps/daemon/internal/localworkspace/native_files.go b/apps/daemon/internal/localworkspace/native_files.go index e669bc66a..7748ca783 100644 --- a/apps/daemon/internal/localworkspace/native_files.go +++ b/apps/daemon/internal/localworkspace/native_files.go @@ -115,12 +115,12 @@ func (b *Binding) writeNativeFile(ctx context.Context, path string, data []byte) } defer root.Close() if err = root.MkdirAll(filepath.Dir(local), 0700); err != nil { - return result, agent.ErrWorkspaceWriteRejected + return result, unpublishedWriteError(err) } temporary := ".oac-write-" + uuid.NewString() file, err := root.OpenFile(temporary, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0600) if err != nil { - return result, agent.ErrWorkspaceWriteRejected + return result, unpublishedWriteError(err) } defer root.Remove(temporary) _, err = file.Write(data) @@ -129,7 +129,7 @@ func (b *Binding) writeNativeFile(ctx context.Context, path string, data []byte) } closeErr := file.Close() if err != nil || closeErr != nil { - return result, agent.ErrWorkspaceWriteRejected + return result, unpublishedWriteError(errors.Join(err, closeErr)) } // Link publishes complete bytes without replacing an existing destination. if err = root.Link(temporary, local); err != nil { @@ -138,12 +138,23 @@ func (b *Binding) writeNativeFile(ctx context.Context, path string, data []byte) return result, agent.ErrWorkspaceWriteDirectory } return result, agent.ErrWorkspaceWriteUnsafe + } else if errors.Is(e, fs.ErrNotExist) { + return result, unpublishedWriteError(err) } return result, agent.ErrWorkspaceWriteRejected } return agent.WorkspaceWriteResult{SizeBytes: int64(len(data))}, nil } +// unpublishedWriteError applies only before this operation publishes a file. +// Parent directories or a temporary file may already have been created. +func unpublishedWriteError(err error) error { + if storageExhausted(err) { + return agent.ErrWorkspaceWriteUnavailable + } + return agent.ErrWorkspaceWriteRejected +} + const artifactFileBytes int64 = 200 << 20 const artifactBatchBytes int64 = 500 << 20 const artifactEntries = 4096 diff --git a/apps/daemon/internal/localworkspace/native_write_error_unix.go b/apps/daemon/internal/localworkspace/native_write_error_unix.go new file mode 100644 index 000000000..6469b1b69 --- /dev/null +++ b/apps/daemon/internal/localworkspace/native_write_error_unix.go @@ -0,0 +1,12 @@ +//go:build linux || darwin + +package localworkspace + +import ( + "errors" + "syscall" +) + +func storageExhausted(err error) bool { + return errors.Is(err, syscall.ENOSPC) || errors.Is(err, syscall.EDQUOT) +} diff --git a/apps/daemon/internal/localworkspace/native_write_error_unix_test.go b/apps/daemon/internal/localworkspace/native_write_error_unix_test.go new file mode 100644 index 000000000..9c4b62370 --- /dev/null +++ b/apps/daemon/internal/localworkspace/native_write_error_unix_test.go @@ -0,0 +1,28 @@ +//go:build linux || darwin + +package localworkspace + +import ( + "errors" + "fmt" + "os" + "syscall" + "testing" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" +) + +func TestNativeUnpublishedWriteStorageFailure(t *testing.T) { + for _, cause := range []error{syscall.ENOSPC, syscall.EDQUOT} { + for _, err := range []error{cause, &os.PathError{Op: "write", Path: "private/path", Err: cause}, fmt.Errorf("sync failed: %w", cause), errors.Join(syscall.EIO, cause)} { + if got := unpublishedWriteError(err); got != agent.ErrWorkspaceWriteUnavailable { + t.Fatalf("storage exhaustion: got %v, want unavailable", got) + } + } + } + for _, cause := range []error{syscall.EIO, syscall.EACCES, syscall.EEXIST, syscall.ENOTDIR} { + if got := unpublishedWriteError(&os.PathError{Op: "write", Path: "private/path", Err: cause}); got != agent.ErrWorkspaceWriteRejected { + t.Fatalf("changed ordinary rejection: %v", got) + } + } +} diff --git a/apps/daemon/internal/localworkspace/native_write_error_windows.go b/apps/daemon/internal/localworkspace/native_write_error_windows.go new file mode 100644 index 000000000..b577072d4 --- /dev/null +++ b/apps/daemon/internal/localworkspace/native_write_error_windows.go @@ -0,0 +1,11 @@ +package localworkspace + +import ( + "errors" + + "golang.org/x/sys/windows" +) + +func storageExhausted(err error) bool { + return errors.Is(err, windows.ERROR_DISK_FULL) || errors.Is(err, windows.ERROR_HANDLE_DISK_FULL) || errors.Is(err, windows.ERROR_DISK_QUOTA_EXCEEDED) +} diff --git a/apps/daemon/internal/localworkspace/native_write_error_windows_test.go b/apps/daemon/internal/localworkspace/native_write_error_windows_test.go new file mode 100644 index 000000000..c250a0cd3 --- /dev/null +++ b/apps/daemon/internal/localworkspace/native_write_error_windows_test.go @@ -0,0 +1,20 @@ +package localworkspace + +import ( + "os" + "testing" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" + "golang.org/x/sys/windows" +) + +func TestNativeUnpublishedWriteStorageFailure(t *testing.T) { + for _, cause := range []error{windows.ERROR_DISK_FULL, windows.ERROR_HANDLE_DISK_FULL, windows.ERROR_DISK_QUOTA_EXCEEDED} { + if got := unpublishedWriteError(&os.PathError{Op: "write", Path: "private/path", Err: cause}); got != agent.ErrWorkspaceWriteUnavailable { + t.Fatalf("storage exhaustion: got %v, want unavailable", got) + } + } + if got := unpublishedWriteError(windows.ERROR_ACCESS_DENIED); got != agent.ErrWorkspaceWriteRejected { + t.Fatalf("changed ordinary rejection: %v", got) + } +} diff --git a/contracts/agents-api/environment-files.md b/contracts/agents-api/environment-files.md index 783da51fe..4f00b04f9 100644 --- a/contracts/agents-api/environment-files.md +++ b/contracts/agents-api/environment-files.md @@ -77,7 +77,7 @@ Unless the table names a code, the 400 errors have type and code `invalid_reques - Before sending any bytes, Core records the write under the Session lock. While input is pending, a Turn is running or an earlier write is unsettled, a new write returns 409 `turn_conflict`. An unsettled write also makes new messages to the Session return 409. - The Runtime checks the complete body against its digest before creating anything, so incomplete input creates nothing. A write that fails later can leave newly created empty parent directories. -- Only a definite receipt from the Runtime settles a write, as committed or rejected. A rejected write changes nothing and releases the Session. If the connection drops, the request times out or no receipt arrives, the request returns 503 and the write stays unsettled, across Core restarts. Core never resends it and has no automatic recovery, so the Session accepts no further writes or messages. Reads still work. +- Only a definite receipt from the Runtime settles a write, as committed or rejected. A rejected write publishes no destination file and releases the Session; newly created empty parent directories can remain. Storage or quota exhaustion confirmed before publication returns 503, while malformed input and destination conflicts retain their existing errors. If the connection drops, the request times out or no receipt arrives, the request returns 503 and the write stays unsettled, across Core restarts. Core never resends it and has no automatic recovery, so the Session accepts no further writes or messages. Reads still work. - Deleting the source File after its bytes were read does not affect the copy. ## Artifacts diff --git a/contracts/agents-api/zh/environment-files.md b/contracts/agents-api/zh/environment-files.md index 1977f2b63..8f31c84e2 100644 --- a/contracts/agents-api/zh/environment-files.md +++ b/contracts/agents-api/zh/environment-files.md @@ -1,7 +1,7 @@ --- title: "Environment 文件与 Artifact" source: contracts/agents-api/environment-files.md -source_hash: 1b58aa02aaccddb9675ef41ebfe2506a6fba0bb12139efb67e0da378d879aee7 +source_hash: f87f06138789666b91140c15ffd104cffba4566580921b94b094bf31591d201d --- Session 工作区保存由 agent 及其工具修改的实时文件。`/agents/environments/{environment_id}/files` 列出一个工作区目录,并在其中创建文件。Turn 完成时,Core 将工作区 `outputs/` 目录中的文件复制为不可变 Artifact,通过 `/agents/sessions/{session_id}/artifacts` 读取。Artifact 的生命周期长于 Environment;工作区文件则不是。 @@ -79,7 +79,7 @@ Session 工作区保存由 agent 及其工具修改的实时文件。`/agents/en - 发送任何字节之前,Core 在 Session 锁下记录写入。输入待处理、Turn 运行或较早写入未结算时,新写入返回 409 `turn_conflict`。未结算写入也使 Session 新消息返回 409。 - Runtime 在创建任何内容前根据摘要检查完整正文,因此不完整输入不创建内容。后续失败的写入可能留下新建空父目录。 -- 仅 Runtime 的确定回执将写入结算为已提交或已拒绝。被拒绝写入不改变内容并释放 Session。连接断开、请求超时或无回执时,返回 503,写入持续未结算,Core 重启后仍如此。Core 不重发,也无自动恢复,因此 Session 不再接受写入或消息。读取仍可用。 +- 仅 Runtime 的确定回执将写入结算为已提交或已拒绝。被拒绝写入不会发布目标文件,并释放 Session;新建的空父目录可能保留。确认在发布前发生的存储或配额耗尽返回 503,格式错误的输入和目标冲突仍保持原有错误。连接断开、请求超时或无回执时,返回 503,写入持续未结算,Core 重启后仍如此。Core 不重发,也无自动恢复,因此 Session 不再接受写入或消息。读取仍可用。 - 源 File 字节读取后,删除该 File 不影响副本。 ## Artifact {#artifacts} diff --git a/services/core/tests/integration/environment_file_write_semantics_public_test.go b/services/core/tests/integration/environment_file_write_semantics_public_test.go index eef998402..c8e66c22e 100644 --- a/services/core/tests/integration/environment_file_write_semantics_public_test.go +++ b/services/core/tests/integration/environment_file_write_semantics_public_test.go @@ -59,7 +59,7 @@ func TestEnvironmentFileCreateRejectionsLeaveNoReceiptOrConsumption(t *testing.T {OrganizationID: "test-org", ProjectID: uuid.NewString(), SubjectKind: "service_account", SubjectID: "test-runner", TokenSHA256: runtimedevice.HashCredential(token), TenantID: h.tenant}, {OrganizationID: "test-org", ProjectID: uuid.NewString(), SubjectKind: "service_account", SubjectID: "tenant-b", TokenSHA256: runtimedevice.HashCredential(other), TenantID: uuid.NewString()}, }) - handler, err := publicHandler(t, h.s, auth, "codex", workerExecution(t, w)) + handler, err := publicHandler(t, h.s, auth, "codex", workerExecution(t, w), acceptUnavailable(t)) if err != nil { t.Fatal(err) } @@ -144,6 +144,16 @@ func TestEnvironmentFileCreateRejectionsLeaveNoReceiptOrConsumption(t *testing.T t.Fatal("rejection did not settle", intent, err) } } + // Exhaustion rejected before publication is unavailable, not invalid input; + // its exact receipt still settles ownership so another write can proceed. + doneUnavailable := post(token, environment.ID, inline) + unavailableID := serveFileWrite(h, proto.WorkspaceWriteResultPayload{Outcome: "rejected", ErrorCode: "resource_unavailable"}) + if got := await(doneUnavailable); got.status != http.StatusServiceUnavailable { + t.Fatal("resource rejection became invalid input", got) + } + if intent, err := FixtureFileWrite(t.Context(), h.s.pool, h.tenant, environment.ID, unavailableID); err != nil || intent.State != "rejected" { + t.Fatal("known unavailable result retained uncertainty", intent, err) + } // The rejected copy did not consume its Source File. if got, err := fileStore.Get(t.Context(), h.tenant, source.ID); err != nil || got.SizeBytes != 3 { t.Fatal("source file consumed", got, err) @@ -156,7 +166,7 @@ func TestEnvironmentFileCreateRejectionsLeaveNoReceiptOrConsumption(t *testing.T if foreign.status != 404 || foreign != missing { t.Fatal("tenant B reached the Environment", foreign, missing) } - if got := states(); !reflect.DeepEqual(got, map[string]int{"rejected": 3}) { + if got := states(); !reflect.DeepEqual(got, map[string]int{"rejected": 4}) { t.Fatal("rejections left a receipt or blocking intent", got) } // Known rejections release the mutation owner for a successor. @@ -165,7 +175,7 @@ func TestEnvironmentFileCreateRejectionsLeaveNoReceiptOrConsumption(t *testing.T if got := await(done); got.status != 201 { t.Fatal("successor rejected", got) } - if got := states(); !reflect.DeepEqual(got, map[string]int{"rejected": 3, "committed": 1}) { + if got := states(); !reflect.DeepEqual(got, map[string]int{"rejected": 4, "committed": 1}) { t.Fatal("successor receipt", got) } session, err := sessionAdapter(h.s).GetSession(t.Context(), h.tenant, h.session.ID)