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
17 changes: 14 additions & 3 deletions apps/daemon/internal/localworkspace/native_files.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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 {
Expand All @@ -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
Expand Down
12 changes: 12 additions & 0 deletions apps/daemon/internal/localworkspace/native_write_error_unix.go
Original file line number Diff line number Diff line change
@@ -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)
}
Original file line number Diff line number Diff line change
@@ -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)
}
}
}
11 changes: 11 additions & 0 deletions apps/daemon/internal/localworkspace/native_write_error_windows.go
Original file line number Diff line number Diff line change
@@ -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)
}
Original file line number Diff line number Diff line change
@@ -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)
}
}
2 changes: 1 addition & 1 deletion contracts/agents-api/environment-files.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions contracts/agents-api/zh/environment-files.md
Original file line number Diff line number Diff line change
@@ -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;工作区文件则不是。
Expand Down Expand Up @@ -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}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
Expand All @@ -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.
Expand All @@ -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)
Expand Down
Loading