diff --git a/apps/daemon/internal/agent/claudesdk/bridge_output.go b/apps/daemon/internal/agent/claudesdk/bridge_output.go deleted file mode 100644 index e80cd22fe..000000000 --- a/apps/daemon/internal/agent/claudesdk/bridge_output.go +++ /dev/null @@ -1,32 +0,0 @@ -package claudesdk - -import "bufio" - -type bridgeOutput struct { - frames chan []byte - current []byte - err error -} - -func (s *session) bridgeOutput() *bridgeOutput { - output := &bridgeOutput{frames: make(chan []byte, 1)} - go func() { - defer close(output.frames) - scanner := bufio.NewScanner(s.process.Stdout) - scanner.Buffer(make([]byte, 64*1024), 2*1024*1024) - for scanner.Scan() { - raw := append([]byte(nil), scanner.Bytes()...) - if !s.receiveWorkspaceDirectory(raw) { - output.frames <- raw - } - } - output.err = scanner.Err() - if output.err != nil { - s.process.Cancel() - } - }() - return output -} -func (o *bridgeOutput) Scan() bool { var ok bool; o.current, ok = <-o.frames; return ok } -func (o *bridgeOutput) Bytes() []byte { return o.current } -func (o *bridgeOutput) Err() error { return o.err } diff --git a/apps/daemon/internal/agent/claudesdk/contracts.go b/apps/daemon/internal/agent/claudesdk/contracts.go index aaf9e92aa..e425fc638 100644 --- a/apps/daemon/internal/agent/claudesdk/contracts.go +++ b/apps/daemon/internal/agent/claudesdk/contracts.go @@ -5,14 +5,10 @@ import "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" // Every public Harness implements each small contract explicitly. Unsupported // extensions return agent.ErrUnsupportedOperation without native effects. var ( - _ agent.Executor = (*executor)(nil) - _ agent.Turn = (*session)(nil) - _ agent.Session = (*session)(nil) - _ agent.DurableSteerer = (*session)(nil) - _ agent.Steerer = (*session)(nil) - _ agent.FunctionResultSubmitter = (*session)(nil) - _ agent.WorkspaceDirectoryLister = (*session)(nil) - _ agent.WorkspaceWriter = (*session)(nil) - _ agent.WorkspaceDirectoryLister = (*executor)(nil) - _ agent.WorkspaceWriter = (*executor)(nil) + _ agent.Executor = (*executor)(nil) + _ agent.Turn = (*session)(nil) + _ agent.Session = (*session)(nil) + _ agent.DurableSteerer = (*session)(nil) + _ agent.Steerer = (*session)(nil) + _ agent.FunctionResultSubmitter = (*session)(nil) ) diff --git a/apps/daemon/internal/agent/claudesdk/declaration.go b/apps/daemon/internal/agent/claudesdk/declaration.go index c17df1372..04deb951c 100644 --- a/apps/daemon/internal/agent/claudesdk/declaration.go +++ b/apps/daemon/internal/agent/claudesdk/declaration.go @@ -116,7 +116,7 @@ func discoverWithCheck(parent context.Context, options agent.DiscoveryOptions, d } caps := &out.Info.Capabilities caps.EnvironmentNone, caps.FunctionTools = proto.CapabilityUnsupported, proto.CapabilityFromBool(info.SupportsWorkspaceFunctions()) - caps.LocalEnvironment, caps.WorkspaceReadPreparation = proto.CapabilitySupported, proto.CapabilitySupported + caps.LocalEnvironment = proto.CapabilitySupported caps.NativeSessionRecovery = proto.CapabilitySupported } out.Info.Available, out.Info.Version = true, info.SDK diff --git a/apps/daemon/internal/agent/claudesdk/declaration_test.go b/apps/daemon/internal/agent/claudesdk/declaration_test.go index 9181b5454..9bd2721ba 100644 --- a/apps/daemon/internal/agent/claudesdk/declaration_test.go +++ b/apps/daemon/internal/agent/claudesdk/declaration_test.go @@ -116,7 +116,7 @@ func TestRuntimeDiscoveryConfigurationAndRegistration(t *testing.T) { } continue } - if calls != 1 || runtime.Info.Available != ready || (runtime.Executor != nil) != ready || runtime.Info.Capabilities.WorkspaceReadPreparation.IsSupported() { + if calls != 1 || runtime.Info.Available != ready || (runtime.Executor != nil) != ready || runtime.Info.Capabilities.LocalEnvironment.IsSupported() { t.Fatalf("runtime: %+v", runtime) } registry := agent.NewRegistry() diff --git a/apps/daemon/internal/agent/claudesdk/executor.go b/apps/daemon/internal/agent/claudesdk/executor.go index 31934e1f6..6c68f9387 100644 --- a/apps/daemon/internal/agent/claudesdk/executor.go +++ b/apps/daemon/internal/agent/claudesdk/executor.go @@ -1,6 +1,7 @@ package claudesdk import ( + "bufio" "context" "encoding/json" "errors" @@ -69,7 +70,6 @@ func startExecutor(ctx context.Context, checked *runtimeCheckCache, probe Config if err != nil { return nil, err } - base.directories.supported = slices.Contains(info.Features, "workspace_directory") e := &executor{base: base, start: start, ready: make(chan error, 1), done: make(chan struct{}), nativeID: start.Resume} go e.read() if err = e.write(start); err == nil { @@ -136,7 +136,8 @@ func (e *executor) write(value any) error { func (e *executor) read() { stderrDone := make(chan struct{}) go func() { _, _ = io.Copy(io.Discard, e.base.process.Stderr); close(stderrDone) }() - scanner := e.base.bridgeOutput() + scanner := bufio.NewScanner(e.base.process.Stdout) + scanner.Buffer(make([]byte, 64*1024), 2*1024*1024) ready := false readyFailure := errors.New("claudesdk: Executor readiness failed") for scanner.Scan() { @@ -179,7 +180,6 @@ func (e *executor) read() { } <-stderrDone _ = e.base.process.Wait() - e.base.stopWorkspaceDirectories() e.mu.Lock() e.invalid = true if e.active != nil { @@ -317,7 +317,3 @@ func (s *session) invalidate() { s.process.Cancel() } } - -func (e *executor) ListWorkspaceDirectory(ctx context.Context, path string, maxEntries int) (agent.WorkspaceDirectoryResult, error) { - return e.base.ListWorkspaceDirectory(ctx, path, maxEntries) -} diff --git a/apps/daemon/internal/agent/claudesdk/preparation_fixture_test.go b/apps/daemon/internal/agent/claudesdk/preparation_fixture_test.go index b16f8ad93..314aa5dbc 100644 --- a/apps/daemon/internal/agent/claudesdk/preparation_fixture_test.go +++ b/apps/daemon/internal/agent/claudesdk/preparation_fixture_test.go @@ -38,7 +38,7 @@ func preparationRequest() proto.PromptRequestPayload { func runPreparationHelper() { mode := os.Getenv("SDK_HELPER_MODE") if len(os.Args) > 1 && strings.HasSuffix(os.Args[1], "runtime_check.js") { - features := []string{"workspace_tools", "workspace_prepare", "workspace_command_observations", "workspace_directory"} + features := []string{"workspace_tools", "workspace_prepare", "workspace_command_observations"} if mode == "old-runtime" { features = []string{"workspace_tools"} } else if mode == "old-command-runtime" { @@ -85,10 +85,6 @@ func runPreparationHelper() { } } emit(bridgeEvent{Type: "executor_ready", Protocol: 3}) - if strings.HasPrefix(mode, "directory-") { - runWorkspaceDirectoryHelper(scanner, state, mode) - return - } if !scanner.Scan() { return } diff --git a/apps/daemon/internal/agent/claudesdk/session.go b/apps/daemon/internal/agent/claudesdk/session.go index fcf91d376..7cfc8c59a 100644 --- a/apps/daemon/internal/agent/claudesdk/session.go +++ b/apps/daemon/internal/agent/claudesdk/session.go @@ -23,13 +23,12 @@ type session struct { cancelOutput chan struct{} nativeEnded bool - directories workspaceDirectoryState - process *clirunner.Process - writeMu *sync.Mutex - functions functionState - steering steeringState - settled chan struct{} - outcome proto.DonePayload + process *clirunner.Process + writeMu *sync.Mutex + functions functionState + steering steeringState + settled chan struct{} + outcome proto.DonePayload } type bridgeEvent struct { diff --git a/apps/daemon/internal/agent/claudesdk/unsupported.go b/apps/daemon/internal/agent/claudesdk/unsupported.go deleted file mode 100644 index 154d904a4..000000000 --- a/apps/daemon/internal/agent/claudesdk/unsupported.go +++ /dev/null @@ -1,15 +0,0 @@ -package claudesdk - -import ( - "context" - - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" -) - -func (s *executor) WriteWorkspaceFile(context.Context, string, []byte) (agent.WorkspaceWriteResult, error) { - return agent.WorkspaceWriteResult{}, agent.ErrWorkspaceWriteUnsupported -} - -func (s *session) WriteWorkspaceFile(context.Context, string, []byte) (agent.WorkspaceWriteResult, error) { - return agent.WorkspaceWriteResult{}, agent.ErrWorkspaceWriteUnsupported -} diff --git a/apps/daemon/internal/agent/claudesdk/unsupported_test.go b/apps/daemon/internal/agent/claudesdk/unsupported_test.go deleted file mode 100644 index 7fe6048c5..000000000 --- a/apps/daemon/internal/agent/claudesdk/unsupported_test.go +++ /dev/null @@ -1,34 +0,0 @@ -package claudesdk - -import ( - "context" - "errors" - "strings" - "testing" - - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" -) - -// Nil concrete receivers prove Unsupported needs no native owner, transport, -// workspace access or input retention. Callers still preserve the separate -// nil-Turn ownership rule on required lifecycle methods. -func TestUnsupportedExtensionsHaveNoNativeEffects(t *testing.T) { - ctx := context.Background() - const secret = "private-fixture-value" - check := func(err error) { - t.Helper() - if !errors.Is(err, agent.ErrUnsupportedOperation) || err.Error() == agent.ErrUnsupportedOperation.Error() { - t.Fatalf("expected explicit unsupported result with reason, got %v", err) - } - if strings.Contains(err.Error(), secret) { - t.Fatal("unsupported error disclosed input") - } - } - for _, owner := range []agent.WorkspaceWriter{(*executor)(nil), (*session)(nil)} { - result, err := owner.WriteWorkspaceFile(ctx, secret, []byte(secret)) - check(err) - if result.SizeBytes != 0 { - t.Fatal("unsupported write fabricated receipt") - } - } -} diff --git a/apps/daemon/internal/agent/claudesdk/workspace_directory.go b/apps/daemon/internal/agent/claudesdk/workspace_directory.go deleted file mode 100644 index 1bde6cbe5..000000000 --- a/apps/daemon/internal/agent/claudesdk/workspace_directory.go +++ /dev/null @@ -1,237 +0,0 @@ -package claudesdk - -import ( - "bytes" - "context" - "encoding/json" - "io" - "io/fs" - "strings" - "sync" - "time" - - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" - "github.com/google/uuid" -) - -const workspaceDirectoryMaxEntries = 1000 -const workspaceDirectoryTimeout = 12 * time.Second - -type workspaceDirectoryState struct { - mu sync.Mutex - supported bool - closed bool - uncertain bool - pending *workspaceDirectory -} -type workspaceDirectory struct { - id string - maxEntries int - done chan struct{} - result agent.WorkspaceDirectoryResult - err error -} -type workspaceDirectoryEvent struct { - Type string `json:"type"` - ID string `json:"id"` - Entries *[]workspaceDirectoryEntry `json:"entries"` - Truncated *bool `json:"truncated"` - Error string `json:"error"` -} - -var _ agent.WorkspaceDirectoryLister = (*session)(nil) - -func (s *session) ListWorkspaceDirectory(ctx context.Context, path string, maxEntries int) (agent.WorkspaceDirectoryResult, error) { - if s.owner != nil { - return s.owner.base.ListWorkspaceDirectory(ctx, path, maxEntries) - } - read, err := s.admitWorkspaceDirectory(ctx, path, maxEntries) - if err != nil { - return agent.WorkspaceDirectoryResult{}, err - } - return s.awaitWorkspaceDirectory(ctx, read) -} - -func (s *session) admitWorkspaceDirectory(ctx context.Context, path string, maxEntries int) (*workspaceDirectory, error) { - w := &s.directories - w.mu.Lock() - defer w.mu.Unlock() - if !w.supported { - return nil, agent.ErrWorkspaceReadUnsupported - } - if w.uncertain { - return nil, agent.ErrWorkspaceReadUncertain - } - if w.closed || ctx == nil || ctx.Err() != nil || s.process.Context().Err() != nil { - return nil, agent.ErrWorkspaceReadUnavailable - } - if maxEntries < 1 || maxEntries > workspaceDirectoryMaxEntries || len(path) > 8192 || strings.ContainsAny(path, "\x00\\\r\n") { - return nil, agent.ErrWorkspaceReadInvalid - } - for _, part := range strings.Split(path, "/") { - if path == "" { - break - } - if part == "" || part == "." || part == ".." { - return nil, agent.ErrWorkspaceReadInvalid - } - } - if w.pending != nil { - return nil, agent.ErrWorkspaceReadBusy - } - read := &workspaceDirectory{id: uuid.NewString(), maxEntries: maxEntries, done: make(chan struct{})} - frame, err := json.Marshal(struct { - Type string `json:"type"` - ID string `json:"id"` - Path string `json:"directory"` - MaxEntries int `json:"max_entries"` - }{"workspace_directory", read.id, path, maxEntries}) - if err != nil || len(frame)+1 > 8192 { - return nil, agent.ErrWorkspaceReadInvalid - } - w.pending = read - deadline := time.Now().Add(workspaceDirectoryTimeout) - if requested, ok := ctx.Deadline(); ok && requested.Before(deadline) { - deadline = requested - } - go func() { - timer := time.NewTimer(time.Until(deadline)) - defer timer.Stop() - select { - case <-read.done: - return - case <-timer.C: - } - s.failWorkspaceDirectory(read) - }() - go func() { - s.writeMu.Lock() - _, err := s.process.Stdin.Write(append(frame, '\n')) - s.writeMu.Unlock() - if err != nil { - s.failWorkspaceDirectory(read) - } - }() - return read, nil -} - -func (s *session) failWorkspaceDirectory(read *workspaceDirectory) { - s.directories.mu.Lock() - pending := s.directories.pending == read - if pending { - s.directories.uncertain = true - s.directories.pending = nil - read.err = agent.ErrWorkspaceReadUncertain - s.process.Cancel() - close(read.done) - } - s.directories.mu.Unlock() -} - -func (s *session) awaitWorkspaceDirectory(_ context.Context, read *workspaceDirectory) (agent.WorkspaceDirectoryResult, error) { - // Caller cancellation cannot discard an already admitted native wait. - <-read.done - return read.result, read.err -} - -func (s *session) receiveWorkspaceDirectory(raw []byte) bool { - var event workspaceDirectoryEvent - if json.Unmarshal(raw, &event) != nil || event.Type != "workspace_directory" { - return false - } - w := &s.directories - w.mu.Lock() - defer w.mu.Unlock() - read := w.pending - if read == nil || event.ID != read.id { - w.uncertain = true - s.process.Cancel() - return true - } - read.err = agent.ErrWorkspaceReadUncertain - decoder := json.NewDecoder(bytes.NewReader(raw)) - decoder.DisallowUnknownFields() - valid := decoder.Decode(&event) == nil && decoder.Decode(new(any)) == io.EOF - if !valid { - w.uncertain = true - w.pending = nil - close(read.done) - s.process.Cancel() - return true - } - if event.Error != "" && event.Entries == nil && event.Truncated == nil { - switch event.Error { - case "not_found": - read.err = fs.ErrNotExist - case "permission": - read.err = fs.ErrPermission - case "invalid": - read.err = agent.ErrWorkspaceReadInvalid - case "busy": - read.err = agent.ErrWorkspaceReadBusy - case "unavailable": - read.err = agent.ErrWorkspaceReadUnavailable - } - } else if event.Error == "" && event.Entries != nil && event.Truncated != nil { - entries, valid := validWorkspaceDirectoryEntries(*event.Entries, read.maxEntries) - if valid && (!*event.Truncated || len(entries) == read.maxEntries) { - read.result = agent.WorkspaceDirectoryResult{Entries: entries, Truncated: *event.Truncated} - read.err = nil - } - } - if read.err == agent.ErrWorkspaceReadUncertain { - w.uncertain = true - s.process.Cancel() - } - w.pending = nil - close(read.done) - return true -} - -func (s *session) stopWorkspaceDirectories() { - w := &s.directories - w.mu.Lock() - defer w.mu.Unlock() - w.closed = true - if w.pending != nil { - read := w.pending - w.pending = nil - w.uncertain = true - read.err = agent.ErrWorkspaceReadUncertain - close(read.done) - } -} - -type workspaceDirectoryEntry struct { - Name string `json:"name"` - Kind string `json:"kind"` - SizeBytes *int64 `json:"size_bytes,omitempty"` -} - -func validWorkspaceDirectoryEntries(raw []workspaceDirectoryEntry, maxEntries int) ([]agent.WorkspaceDirectoryEntry, bool) { - if raw == nil || len(raw) > maxEntries { - return nil, false - } - entries := make([]agent.WorkspaceDirectoryEntry, 0, len(raw)) - names := make(map[string]bool, len(raw)) - for _, entry := range raw { - if entry.Name == "" || entry.Name == "." || entry.Name == ".." || strings.ContainsAny(entry.Name, "/\x00") || names[entry.Name] { - return nil, false - } - names[entry.Name] = true - switch entry.Kind { - case "file": - if entry.SizeBytes == nil || *entry.SizeBytes < 0 { - return nil, false - } - case "directory", "symlink", "other": - if entry.SizeBytes != nil { - return nil, false - } - default: - return nil, false - } - entries = append(entries, agent.WorkspaceDirectoryEntry{Name: entry.Name, Kind: entry.Kind, SizeBytes: entry.SizeBytes}) - } - return entries, true -} diff --git a/apps/daemon/internal/agent/claudesdk/workspace_directory_live_linux_test.go b/apps/daemon/internal/agent/claudesdk/workspace_directory_live_linux_test.go deleted file mode 100644 index 03445a78e..000000000 --- a/apps/daemon/internal/agent/claudesdk/workspace_directory_live_linux_test.go +++ /dev/null @@ -1,99 +0,0 @@ -//go:build linux - -package claudesdk - -import ( - "context" - "encoding/json" - "errors" - "io/fs" - "os" - "path/filepath" - "testing" - - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" -) - -func liveWorkspaceDirectoryFixtures(t *testing.T, root string) { - t.Helper() - for _, name := range []string{"directory-empty", "directory-denied"} { - if err := os.Mkdir(filepath.Join(root, name), 0700); err != nil { - t.Fatal(err) - } - } - if err := os.Chmod(filepath.Join(root, "directory-denied"), 0); err != nil { - t.Fatal(err) - } - for name, target := range map[string]string{"directory-inside-link": "directory-empty", "directory-outside-link": filepath.Dir(root)} { - if err := os.Symlink(target, filepath.Join(root, name)); err != nil { - t.Fatal(err) - } - } -} - -func liveWorkspaceDirectories(t *testing.T, ctx context.Context, lister agent.WorkspaceDirectoryLister, root, stage string) { - t.Helper() - result, err := lister.ListWorkspaceDirectory(ctx, "", 1000) - if err != nil || result.Truncated { - t.Fatal("workspace directory failed", stage, err) - } - expected, err := os.ReadDir(root) - if err != nil || len(expected) != len(result.Entries) { - t.Fatal("directory length differs", stage, err) - } - found := make(map[string]agent.WorkspaceDirectoryEntry) - for _, entry := range result.Entries { - found[entry.Name] = entry - } - for _, entry := range expected { - actual, ok := found[entry.Name()] - if !ok { - t.Fatal("missing directory entry", stage, entry.Name()) - } - stat, err := os.Lstat(filepath.Join(root, entry.Name())) - if err != nil { - t.Fatal(err) - } - switch { - case stat.Mode().IsRegular(): - // The command heartbeat is intentionally changing during live execution. - if actual.Kind != "file" || actual.SizeBytes == nil || (entry.Name() != "heartbeat.txt" && *actual.SizeBytes != stat.Size()) { - t.Fatal("file metadata differs", stage, actual) - } - case stat.IsDir(): - if actual.Kind != "directory" || actual.SizeBytes != nil { - t.Fatal(actual) - } - case stat.Mode()&os.ModeSymlink != 0: - if actual.Kind != "symlink" || actual.SizeBytes != nil { - t.Fatal(actual) - } - } - } - empty, err := lister.ListWorkspaceDirectory(ctx, "directory-empty", 2) - if err != nil || empty.Truncated || len(empty.Entries) != 0 { - t.Fatal(empty, err) - } - bounded, err := lister.ListWorkspaceDirectory(ctx, "", 1) - if err != nil || !bounded.Truncated || len(bounded.Entries) != 1 { - t.Fatal(bounded, err) - } - for _, path := range []string{"directory-inside-link", "directory-outside-link"} { - if _, err := lister.ListWorkspaceDirectory(ctx, path, 2); !errors.Is(err, fs.ErrPermission) && !errors.Is(err, agent.ErrWorkspaceReadInvalid) { - t.Fatal("directory denial missing", path, err) - } - } - if _, err := lister.ListWorkspaceDirectory(ctx, "directory-denied", 2); !errors.Is(err, fs.ErrPermission) { - t.Fatal("permission denial missing", err) - } - if _, err := lister.ListWorkspaceDirectory(ctx, "directory-missing", 2); !errors.Is(err, fs.ErrNotExist) { - t.Fatal(err) - } - raw, err := json.MarshalIndent(result, "", " ") - if err != nil { - t.Fatal(err) - } - if err := os.WriteFile(filepath.Join(filepath.Dir(root), "directory-"+stage+".json"), raw, 0600); err != nil { - t.Fatal(err) - } -} diff --git a/apps/daemon/internal/agent/claudesdk/workspace_directory_test.go b/apps/daemon/internal/agent/claudesdk/workspace_directory_test.go deleted file mode 100644 index e073cffae..000000000 --- a/apps/daemon/internal/agent/claudesdk/workspace_directory_test.go +++ /dev/null @@ -1,214 +0,0 @@ -//go:build unix - -package claudesdk - -import ( - "bufio" - - "context" - "encoding/json" - "errors" - "io/fs" - "os" - "path/filepath" - "strings" - "testing" - "time" - - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" - "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" -) - -func runWorkspaceDirectoryHelper(scanner *bufio.Scanner, state, mode string) { - turnID := "" - for scanner.Scan() { - var req struct { - Type, ID string - TurnID string `json:"turn_id"` - Path string `json:"directory"` - Max int `json:"max_entries"` - } - _ = json.Unmarshal(scanner.Bytes(), &req) - if req.Type == "turn_cancel" { - reusable := false - _ = json.NewEncoder(os.Stdout).Encode(bridgeEvent{Type: "error", TurnID: turnID, Code: "cancelled"}) - confirmed := true - _ = json.NewEncoder(os.Stdout).Encode(bridgeEvent{Type: "turn_settled", Confirmed: &confirmed, TurnID: turnID, Reusable: &reusable, Reason: "fixture_closed"}) - return - } - if req.Type == "turn_start" { - turnID = req.TurnID - _ = json.NewEncoder(os.Stdout).Encode(bridgeEvent{Type: "turn_started", TurnID: turnID}) - _ = os.WriteFile(filepath.Join(state, "directory-started"), []byte("{}"), 0600) - continue - } - _ = os.WriteFile(filepath.Join(state, "directory-admitted"), []byte("{}"), 0600) - if mode == "directory-held" { - for { - if _, err := os.Stat(filepath.Join(state, "directory-release")); err == nil { - break - } - time.Sleep(time.Millisecond) - } - } - if mode == "directory-exit" { - return - } - entries := []map[string]any{{"name": "file", "kind": "file", "size_bytes": 3}} - if req.Path == "empty" { - entries = []map[string]any{} - } - response := map[string]any{"type": "workspace_directory", "id": req.ID, "entries": entries, "truncated": false} - switch req.Path { - case "uncertain", "invalid", "not_found", "permission": - response = map[string]any{"type": "workspace_directory", "id": req.ID, "error": req.Path} - case "bad-kind": - entries[0]["kind"] = "unknown" - case "wrong-id": - response["id"] = "other" - case "extra": - response["extra"] = true - case "bad-size": - entries[0]["size_bytes"] = -1 - case "missing-size": - delete(entries[0], "size_bytes") - case "duplicate": - response["entries"] = append(entries, entries[0]) - case "escape-name": - entries[0]["name"] = "../secret" - case "null": - response["entries"] = nil - } - - _ = json.NewEncoder(os.Stdout).Encode(response) - } -} - -func directoryExecutor(t *testing.T, mode string) (*executor, Config) { - t.Helper() - config := preparationFixture(t, mode) - resource, err := NewExecutorFactory(config)(t.Context(), preparationRequest()) - if err != nil { - t.Fatal(err) - } - e := resource.(*executor) - t.Cleanup(func() { _ = e.Close(context.Background()) }) - return e, config -} - -func TestWorkspaceDirectoryBoundsAndMetadata(t *testing.T) { - e, _ := directoryExecutor(t, "directory-normal") - for _, path := range []string{"/absolute", "../escape", "a//b", "a/./b", "a\\b", "a\x00b", strings.Repeat("界", 3000)} { - if _, err := e.ListWorkspaceDirectory(t.Context(), path, 4); !errors.Is(err, agent.ErrWorkspaceReadInvalid) { - t.Fatalf("accepted %q: %v", path, err) - } - } - for _, limit := range []int{0, -1, workspaceDirectoryMaxEntries + 1} { - if _, err := e.ListWorkspaceDirectory(t.Context(), "", limit); err != agent.ErrWorkspaceReadInvalid { - t.Fatal(limit, err) - } - } - result, err := e.ListWorkspaceDirectory(t.Context(), "", 4) - if err != nil || result.Truncated || len(result.Entries) != 1 || result.Entries[0].SizeBytes == nil || *result.Entries[0].SizeBytes != 3 { - t.Fatal(result, err) - } - empty, err := e.ListWorkspaceDirectory(t.Context(), "empty", 4) - if err != nil || len(empty.Entries) != 0 || empty.Entries == nil { - t.Fatal(empty, err) - } - for path, expected := range map[string]error{"not_found": fs.ErrNotExist, "permission": fs.ErrPermission, "invalid": agent.ErrWorkspaceReadInvalid} { - if _, err := e.ListWorkspaceDirectory(t.Context(), path, 4); err != expected { - t.Fatal(path, err) - } - } - if _, err := e.ListWorkspaceDirectory(t.Context(), "", 4); err != nil { - t.Fatal(err) - } -} - -func TestWorkspaceDirectoryDetachAndTurnStart(t *testing.T) { - e, config := directoryExecutor(t, "directory-held") - ctx, detach := context.WithCancel(t.Context()) - defer detach() - result := make(chan error, 1) - go func() { _, err := e.ListWorkspaceDirectory(ctx, "file", 4); result <- err }() - waitPreparationFile(t, filepath.Join(config.StateDir, "directory-admitted")) - detach() - if _, err := e.ListWorkspaceDirectory(t.Context(), "file", 4); err != agent.ErrWorkspaceReadBusy { - t.Fatal(err) - } - out := make(chan proto.Envelope, 16) - running, err := e.StartTurn(t.Context(), "run", proto.TextInput("hello"), out) - if err != nil { - t.Fatal(err) - } - select { - case err := <-result: - t.Fatal("caller detach discarded native wait", err) - default: - } - if err := os.WriteFile(filepath.Join(config.StateDir, "directory-release"), nil, 0600); err != nil { - t.Fatal(err) - } - if err := <-result; err != nil { - t.Fatal(err) - } - waitPreparationFile(t, filepath.Join(config.StateDir, "directory-started")) - if _, err := running.(agent.WorkspaceDirectoryLister).ListWorkspaceDirectory(t.Context(), "file", 4); err != nil { - t.Fatal(err) - } -} - -func TestWorkspaceDirectoryUnknownAndRelease(t *testing.T) { - for _, path := range []string{"uncertain", "wrong-id", "bad-kind", "bad-size", "missing-size", "duplicate", "escape-name", "null", "extra"} { - t.Run(path, func(t *testing.T) { - e, _ := directoryExecutor(t, "directory-normal") - result, err := e.ListWorkspaceDirectory(t.Context(), path, 4) - if err != agent.ErrWorkspaceReadUncertain || len(result.Entries) != 0 { - t.Fatal(result, err) - } - if _, err := e.ListWorkspaceDirectory(t.Context(), "file", 4); err != agent.ErrWorkspaceReadUncertain && err != agent.ErrWorkspaceReadUnavailable { - t.Fatal(err) - } - }) - } - for _, mode := range []string{"directory-held", "directory-exit"} { - t.Run(mode, func(t *testing.T) { - e, config := directoryExecutor(t, mode) - done := make(chan error, 1) - go func() { _, err := e.ListWorkspaceDirectory(t.Context(), "file", 4); done <- err }() - waitPreparationFile(t, filepath.Join(config.StateDir, "directory-admitted")) - if mode == "directory-held" { - if err := e.Close(t.Context()); err != nil { - t.Fatal(err) - } - } - select { - case err := <-done: - if err != agent.ErrWorkspaceReadUncertain { - t.Fatal(err) - } - case <-time.After(5 * time.Second): - t.Fatal("native waiter lost on release") - } - }) - } -} - -func TestWorkspaceDirectoryDeadlineStopsOwnerBeforeUnknown(t *testing.T) { - e, config := directoryExecutor(t, "directory-held") - ctx, cancel := context.WithTimeout(t.Context(), 100*time.Millisecond) - defer cancel() - done := make(chan error, 1) - go func() { _, err := e.ListWorkspaceDirectory(ctx, "file", 4); done <- err }() - waitPreparationFile(t, filepath.Join(config.StateDir, "directory-admitted")) - if err := <-done; err != agent.ErrWorkspaceReadUncertain { - t.Fatal(err) - } - if e.base.process.Context().Err() == nil { - t.Fatal("uncertain deadline returned before owner cancellation") - } - if _, err := e.StartTurn(t.Context(), "late", proto.TextInput("hello"), make(chan proto.Envelope, 8)); err == nil { - t.Fatal("unknown owner accepted a new Turn") - } -} diff --git a/apps/daemon/internal/agent/claudesdk/workspace_live_linux_test.go b/apps/daemon/internal/agent/claudesdk/workspace_live_linux_test.go index f1e2f6fa7..d4b9ab32b 100644 --- a/apps/daemon/internal/agent/claudesdk/workspace_live_linux_test.go +++ b/apps/daemon/internal/agent/claudesdk/workspace_live_linux_test.go @@ -59,7 +59,6 @@ func TestLiveClaudeWorkspaceTurns(t *testing.T) { t.Fatal(err) } } - liveWorkspaceDirectoryFixtures(t, config.Workspace.Directory) heartbeat := filepath.Join(config.Workspace.Directory, "heartbeat.txt") artifact := filepath.Join(config.Workspace.Directory, "value.txt") type evidence struct { @@ -105,7 +104,6 @@ func TestLiveClaudeWorkspaceTurns(t *testing.T) { s := running.(*session) defer s.Cancel(context.Background()) proof.BridgePID = s.process.Cmd.Process.Pid - liveWorkspaceDirectories(t, ctx, s, config.Workspace.Directory, "active") ticker := time.NewTicker(80 * time.Millisecond) defer ticker.Stop() for out != nil { @@ -115,7 +113,6 @@ func TestLiveClaudeWorkspaceTurns(t *testing.T) { case <-ticker.C: value, _ := os.ReadFile(heartbeat) if cancelOnEffect && !proof.Cancelled && len(value) > 0 && string(value) != "0" && string(value) != "1" { - liveWorkspaceDirectories(t, ctx, s, config.Workspace.Directory, "effect") started := time.Now() if err := s.Cancel(ctx); err != nil { t.Fatal("factory cancellation failed", err) diff --git a/apps/daemon/internal/agent/codex/contracts.go b/apps/daemon/internal/agent/codex/contracts.go index a8b58b88e..2514d88f1 100644 --- a/apps/daemon/internal/agent/codex/contracts.go +++ b/apps/daemon/internal/agent/codex/contracts.go @@ -5,14 +5,10 @@ import "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" // Every public Harness implements each small contract explicitly. Unsupported // extensions return agent.ErrUnsupportedOperation without native effects. var ( - _ agent.Executor = (*Executor)(nil) - _ agent.Turn = (*Session)(nil) - _ agent.Session = (*Session)(nil) - _ agent.DurableSteerer = (*Session)(nil) - _ agent.Steerer = (*Session)(nil) - _ agent.FunctionResultSubmitter = (*Session)(nil) - _ agent.WorkspaceDirectoryLister = (*Session)(nil) - _ agent.WorkspaceWriter = (*Session)(nil) - _ agent.WorkspaceDirectoryLister = (*Executor)(nil) - _ agent.WorkspaceWriter = (*Executor)(nil) + _ agent.Executor = (*Executor)(nil) + _ agent.Turn = (*Session)(nil) + _ agent.Session = (*Session)(nil) + _ agent.DurableSteerer = (*Session)(nil) + _ agent.Steerer = (*Session)(nil) + _ agent.FunctionResultSubmitter = (*Session)(nil) ) diff --git a/apps/daemon/internal/agent/codex/declaration.go b/apps/daemon/internal/agent/codex/declaration.go index 2b497b07e..48942f2fc 100644 --- a/apps/daemon/internal/agent/codex/declaration.go +++ b/apps/daemon/internal/agent/codex/declaration.go @@ -61,7 +61,6 @@ func discoverWithCheck(parent context.Context, options agent.DiscoveryOptions, i caps := &runtime.Info.Capabilities caps.NativeSessionRecovery = proto.CapabilityFromBool(SupportsNativeSessionRecovery(version)) caps.LocalEnvironment = proto.CapabilityFromBool(SupportsLocalEnvironment(version)) - caps.WorkspaceReadPreparation = caps.LocalEnvironment caps.MCPHTTPRequired = proto.CapabilityFromBool(SupportsNativeSessionRecovery(version)) runtime.Executor = NewExecutorFactory() runtime.View = discoverView(version) diff --git a/apps/daemon/internal/agent/codex/declaration_test.go b/apps/daemon/internal/agent/codex/declaration_test.go index 65a4460c4..2653897ec 100644 --- a/apps/daemon/internal/agent/codex/declaration_test.go +++ b/apps/daemon/internal/agent/codex/declaration_test.go @@ -15,7 +15,7 @@ func TestMCPRequiredDiscoveryRequiresPinnedNative(t *testing.T) { for _, version := range []string{"codex-cli 0.153.4", "codex-cli 0.153.3", "codex-cli 0.154.0"} { runtime := discoverWithCheck(t.Context(), agent.DiscoveryOptions{Stdout: io.Discard, Stderr: io.Discard}, Declaration.Info, func(context.Context, string) (string, error) { return version, nil }) - if !runtime.Info.Available || runtime.Executor == nil || runtime.Info.Capabilities.WorkspaceReadPreparation.IsSupported() != SupportsLocalEnvironment(version) { + if !runtime.Info.Available || runtime.Executor == nil || runtime.Info.Capabilities.LocalEnvironment.IsSupported() != SupportsLocalEnvironment(version) { t.Fatalf("factories: %+v", runtime) } if runtime.Info.Capabilities.MCPHTTPRequired.IsSupported() != (version == "codex-cli 0.153.4") { diff --git a/apps/daemon/internal/agent/codex/preparation_router_test.go b/apps/daemon/internal/agent/codex/preparation_router_test.go index b743ec312..f343b05e0 100644 --- a/apps/daemon/internal/agent/codex/preparation_router_test.go +++ b/apps/daemon/internal/agent/codex/preparation_router_test.go @@ -68,7 +68,8 @@ func TestPreparationRouterRetainsActualNativeChild(t *testing.T) { return e, nil }) sender := make(preparationWireSender, 64) - r, err := dispatch.New(dispatch.Config{Registry: registry, Sender: sender, LocalWorkspace: binding}) + environments := func(proto.AssignmentRef, proto.AssignmentBindPayload) dispatch.Environment { return binding } + r, err := dispatch.New(dispatch.Config{Registry: registry, Sender: sender, Environments: environments}) if err != nil { t.Fatal(err) } diff --git a/apps/daemon/internal/agent/codex/unsupported.go b/apps/daemon/internal/agent/codex/unsupported.go deleted file mode 100644 index 3688145dc..000000000 --- a/apps/daemon/internal/agent/codex/unsupported.go +++ /dev/null @@ -1,22 +0,0 @@ -package codex - -import ( - "context" - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" -) - -func (s *Executor) ListWorkspaceDirectory(context.Context, string, int) (agent.WorkspaceDirectoryResult, error) { - return agent.WorkspaceDirectoryResult{}, agent.ErrWorkspaceReadUnsupported -} - -func (s *Executor) WriteWorkspaceFile(context.Context, string, []byte) (agent.WorkspaceWriteResult, error) { - return agent.WorkspaceWriteResult{}, agent.ErrWorkspaceWriteUnsupported -} - -func (s *Session) ListWorkspaceDirectory(context.Context, string, int) (agent.WorkspaceDirectoryResult, error) { - return agent.WorkspaceDirectoryResult{}, agent.ErrWorkspaceReadUnsupported -} - -func (s *Session) WriteWorkspaceFile(context.Context, string, []byte) (agent.WorkspaceWriteResult, error) { - return agent.WorkspaceWriteResult{}, agent.ErrWorkspaceWriteUnsupported -} diff --git a/apps/daemon/internal/agent/codex/unsupported_test.go b/apps/daemon/internal/agent/codex/unsupported_test.go deleted file mode 100644 index 8a4234e13..000000000 --- a/apps/daemon/internal/agent/codex/unsupported_test.go +++ /dev/null @@ -1,41 +0,0 @@ -package codex - -import ( - "context" - "errors" - "strings" - "testing" - - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" -) - -// Nil concrete receivers prove Unsupported needs no native owner, transport, -// workspace access or input retention. Callers still preserve the separate -// nil-Turn ownership rule on required lifecycle methods. -func TestUnsupportedExtensionsHaveNoNativeEffects(t *testing.T) { - ctx := context.Background() - const secret = "private-fixture-value" - check := func(err error) { - t.Helper() - if !errors.Is(err, agent.ErrUnsupportedOperation) || err.Error() == agent.ErrUnsupportedOperation.Error() { - t.Fatalf("expected explicit unsupported result with reason, got %v", err) - } - if strings.Contains(err.Error(), secret) { - t.Fatal("unsupported error disclosed input") - } - } - for _, owner := range []agent.WorkspaceDirectoryLister{(*Executor)(nil), (*Session)(nil)} { - result, err := owner.ListWorkspaceDirectory(ctx, secret, 1) - check(err) - if len(result.Entries) != 0 || result.Truncated { - t.Fatal("unsupported listing fabricated entries") - } - } - for _, owner := range []agent.WorkspaceWriter{(*Executor)(nil), (*Session)(nil)} { - result, err := owner.WriteWorkspaceFile(ctx, secret, []byte(secret)) - check(err) - if result.SizeBytes != 0 { - t.Fatal("unsupported write fabricated receipt") - } - } -} diff --git a/apps/daemon/internal/agent/contract_declarations_test.go b/apps/daemon/internal/agent/contract_declarations_test.go index 4479f26d4..fe8ec74e7 100644 --- a/apps/daemon/internal/agent/contract_declarations_test.go +++ b/apps/daemon/internal/agent/contract_declarations_test.go @@ -16,14 +16,12 @@ import ( // decision. Each adapter's compile assertions then enforce its actual methods. func TestPublicHarnessContractDeclarations(t *testing.T) { roles := map[string][]string{ - "Executor": {"executor"}, - "Turn": {"session"}, - "Session": {"session"}, - "DurableSteerer": {"session"}, - "Steerer": {"session"}, - "FunctionResultSubmitter": {"session"}, - "WorkspaceDirectoryLister": {"executor", "session"}, - "WorkspaceWriter": {"executor", "session"}, + "Executor": {"executor"}, + "Turn": {"session"}, + "Session": {"session"}, + "DurableSteerer": {"session"}, + "Steerer": {"session"}, + "FunctionResultSubmitter": {"session"}, } files, err := filepath.Glob("*.go") if err != nil { diff --git a/apps/daemon/internal/agent/harness.go b/apps/daemon/internal/agent/harness.go index e3d0b9826..712169023 100644 --- a/apps/daemon/internal/agent/harness.go +++ b/apps/daemon/internal/agent/harness.go @@ -15,7 +15,8 @@ // static declaration list and installs each resulting Runtime through Register. // Availability and factory selection belong to the adapter. RegisterKind resets // the factories, so Register installs it first. RegisterExecutor derives the -// Preparation capability; the adapter declares WorkspaceReadPreparation. +// Preparation capability. The Runtime's Environment owner, not the adapter, +// serves and declares workspace operations. // // Runtime registration and Core service qualification remain separate. A public // Harness also needs a profile in services/core/internal/engine; advertising @@ -604,30 +605,6 @@ type FunctionResultSubmitter interface { SubmitFunctionResult(context.Context, proto.FunctionResultPayload) error } -// Workspace extensions, implemented by each owner explicitly. The common -// Runtime may supply an authorized workspace owner independently of the adapter. -// A resource without native access returns the corresponding Unsupported error; -// this does not disable capabilities provided by the common workspace owner. - -// WorkspaceDirectoryLister reads one directory through an existing workspace owner. -// An empty directory selects the root; other paths contain only relative components. -// Entry names are single components. Kind is file, directory, symlink or other; -// SizeBytes is present and nonnegative only for regular files. Results have no -// prescribed order, and Truncated must not be presented as a complete inventory. -// maxEntries is positive; adapters may reject limits above their private bound. -// Successful return requires settled directory/metadata access and handle cleanup. -// Implementations retain workspace authorization and isolation and use the existing -// WorkspaceRead errors for unsupported, unavailable, busy, invalid or uncertain reads. -// This interface does not establish public Files pagination or feature admission. -type WorkspaceDirectoryLister interface { - ListWorkspaceDirectory(context.Context, string, int) (WorkspaceDirectoryResult, error) -} - -// WorkspaceWriter confirms a native commit on an already authorized prepared owner. -type WorkspaceWriter interface { - WriteWorkspaceFile(context.Context, string, []byte) (WorkspaceWriteResult, error) -} - // Kind registration. // RegisterKind installs the heartbeat descriptor and model configuration for an diff --git a/apps/daemon/internal/agent/mcode/contracts.go b/apps/daemon/internal/agent/mcode/contracts.go index f58d8ca5e..500870957 100644 --- a/apps/daemon/internal/agent/mcode/contracts.go +++ b/apps/daemon/internal/agent/mcode/contracts.go @@ -5,14 +5,10 @@ import "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" // Every public Harness implements each small contract explicitly. Unsupported // extensions return agent.ErrUnsupportedOperation without native effects. var ( - _ agent.Executor = (*executor)(nil) - _ agent.Turn = (*Session)(nil) - _ agent.Session = (*Session)(nil) - _ agent.DurableSteerer = (*Session)(nil) - _ agent.Steerer = (*Session)(nil) - _ agent.FunctionResultSubmitter = (*Session)(nil) - _ agent.WorkspaceDirectoryLister = (*Session)(nil) - _ agent.WorkspaceWriter = (*Session)(nil) - _ agent.WorkspaceDirectoryLister = (*executor)(nil) - _ agent.WorkspaceWriter = (*executor)(nil) + _ agent.Executor = (*executor)(nil) + _ agent.Turn = (*Session)(nil) + _ agent.Session = (*Session)(nil) + _ agent.DurableSteerer = (*Session)(nil) + _ agent.Steerer = (*Session)(nil) + _ agent.FunctionResultSubmitter = (*Session)(nil) ) diff --git a/apps/daemon/internal/agent/mcode/declaration_test.go b/apps/daemon/internal/agent/mcode/declaration_test.go index 386c4df2d..a2a72c6f6 100644 --- a/apps/daemon/internal/agent/mcode/declaration_test.go +++ b/apps/daemon/internal/agent/mcode/declaration_test.go @@ -21,7 +21,7 @@ func TestMCodeExecutionFollowsAvailability(t *testing.T) { rc := agent.DiscoveryOptions{Stdout: io.Discard, Stderr: io.Discard} runtime := discoverWithCheck(t.Context(), rc, Declaration.Info, func(context.Context, string) (string, error) { return tc.version, tc.check }) info, available := runtime.Info, tc.check == nil - if (runtime.Executor != nil) != available || info.Capabilities.WorkspaceReadPreparation.IsSupported() { + if (runtime.Executor != nil) != available { t.Fatalf("factories: %+v", runtime) } if info.Available != available || info.Capabilities.EnvironmentNone.IsSupported() != available || info.Capabilities.DurableInputReceipts.IsSupported() != available || info.Capabilities.SubagentObservations.IsSupported() != available { diff --git a/apps/daemon/internal/agent/mcode/discovery_workspace.go b/apps/daemon/internal/agent/mcode/discovery_workspace.go index 6ee8d78e7..7c177ec92 100644 --- a/apps/daemon/internal/agent/mcode/discovery_workspace.go +++ b/apps/daemon/internal/agent/mcode/discovery_workspace.go @@ -48,7 +48,7 @@ func discoverWorkspace(parent context.Context, options agent.DiscoveryOptions, r caps := &runtime.Info.Capabilities caps.EnvironmentNone = proto.CapabilityUnsupported - caps.LocalEnvironment, caps.WorkspaceReadPreparation = proto.CapabilitySupported, proto.CapabilitySupported + caps.LocalEnvironment = proto.CapabilitySupported return &c } diff --git a/apps/daemon/internal/agent/mcode/unsupported.go b/apps/daemon/internal/agent/mcode/unsupported.go index 2ee04cf90..85a126b50 100644 --- a/apps/daemon/internal/agent/mcode/unsupported.go +++ b/apps/daemon/internal/agent/mcode/unsupported.go @@ -10,19 +10,3 @@ import ( func (s *Session) SubmitFunctionResult(context.Context, proto.FunctionResultPayload) error { return fmt.Errorf("%w: native public function tools are not qualified", agent.ErrUnsupportedOperation) } - -func (s *executor) ListWorkspaceDirectory(context.Context, string, int) (agent.WorkspaceDirectoryResult, error) { - return agent.WorkspaceDirectoryResult{}, agent.ErrWorkspaceReadUnsupported -} - -func (s *executor) WriteWorkspaceFile(context.Context, string, []byte) (agent.WorkspaceWriteResult, error) { - return agent.WorkspaceWriteResult{}, agent.ErrWorkspaceWriteUnsupported -} - -func (s *Session) ListWorkspaceDirectory(context.Context, string, int) (agent.WorkspaceDirectoryResult, error) { - return agent.WorkspaceDirectoryResult{}, agent.ErrWorkspaceReadUnsupported -} - -func (s *Session) WriteWorkspaceFile(context.Context, string, []byte) (agent.WorkspaceWriteResult, error) { - return agent.WorkspaceWriteResult{}, agent.ErrWorkspaceWriteUnsupported -} diff --git a/apps/daemon/internal/agent/mcode/unsupported_test.go b/apps/daemon/internal/agent/mcode/unsupported_test.go index 3e23d3c24..4d0c48720 100644 --- a/apps/daemon/internal/agent/mcode/unsupported_test.go +++ b/apps/daemon/internal/agent/mcode/unsupported_test.go @@ -27,18 +27,4 @@ func TestUnsupportedExtensionsHaveNoNativeEffects(t *testing.T) { } var turn *Session check(turn.SubmitFunctionResult(ctx, proto.FunctionResultPayload{CallID: secret, DeliveryID: secret})) - for _, owner := range []agent.WorkspaceDirectoryLister{(*executor)(nil), (*Session)(nil)} { - result, err := owner.ListWorkspaceDirectory(ctx, secret, 1) - check(err) - if len(result.Entries) != 0 || result.Truncated { - t.Fatal("unsupported listing fabricated entries") - } - } - for _, owner := range []agent.WorkspaceWriter{(*executor)(nil), (*Session)(nil)} { - result, err := owner.WriteWorkspaceFile(ctx, secret, []byte(secret)) - check(err) - if result.SizeBytes != 0 { - t.Fatal("unsupported write fabricated receipt") - } - } } diff --git a/apps/daemon/internal/agent/workspace_read.go b/apps/daemon/internal/agent/workspace_read.go index 7d4ebbe3b..b156d6c2d 100644 --- a/apps/daemon/internal/agent/workspace_read.go +++ b/apps/daemon/internal/agent/workspace_read.go @@ -1,14 +1,9 @@ package agent -import ( - "errors" - "fmt" -) +import "errors" var ( - ErrWorkspaceReadUnsupported = fmt.Errorf("%w: workspace read", ErrUnsupportedOperation) ErrWorkspaceReadUnavailable = errors.New("workspace read unavailable") - ErrWorkspaceReadBusy = errors.New("workspace read busy") ErrWorkspaceReadInvalid = errors.New("workspace read invalid") ErrWorkspaceReadUncertain = errors.New("workspace read outcome uncertain") // ErrWorkspaceNotDirectory reports that a directory request's own path is diff --git a/apps/daemon/internal/agent/workspace_write.go b/apps/daemon/internal/agent/workspace_write.go index e584843f8..c50b075bd 100644 --- a/apps/daemon/internal/agent/workspace_write.go +++ b/apps/daemon/internal/agent/workspace_write.go @@ -10,7 +10,6 @@ type WorkspaceWriteResult struct { } var ( - ErrWorkspaceWriteUnsupported = fmt.Errorf("%w: workspace write", ErrUnsupportedOperation) ErrWorkspaceWriteUnavailable = errors.New("workspace write unavailable") ErrWorkspaceWriteBusy = errors.New("workspace write busy") ErrWorkspaceWriteInvalid = errors.New("workspace write invalid") diff --git a/apps/daemon/internal/cli/connect.go b/apps/daemon/internal/cli/connect.go index afdea1a96..45b6b79d1 100644 --- a/apps/daemon/internal/cli/connect.go +++ b/apps/daemon/internal/cli/connect.go @@ -315,6 +315,33 @@ func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof } } +// localEnvironments resolves the dedicated local workspace as the Environment +// owner of the one Session it serves. +func localEnvironments(local *localworkspace.Binding) func(proto.AssignmentRef, proto.AssignmentBindPayload) dispatch.Environment { + if local == nil { + return nil + } + return func(ref proto.AssignmentRef, bind proto.AssignmentBindPayload) dispatch.Environment { + if !local.Matches(bind.EnvironmentID, ref.SessionID) { + return nil + } + return local + } +} + +// localEnvironmentKinds declares, for each kind that supports a local +// Environment, the read-only preparation and output export that the local +// workspace owner serves. +func localEnvironmentKinds(registry *agent.Registry, local *localworkspace.Binding) []proto.SupportedAgentKind { + kinds := registry.SupportedAgentKinds() + for i := range kinds { + caps := &kinds[i].Capabilities + caps.WorkspaceReadPreparation = proto.CapabilityFromBool(local != nil && caps.LocalEnvironment.IsSupported()) + caps.WorkspaceOutputExport = caps.WorkspaceReadPreparation + } + return kinds +} + // pumpConn runs the per-connection workload: a dispatch.Router fed by // conn.Recv(), heartbeats every boot.HeartbeatInterval(), and a // confirmed router.Shutdown before returning ownership to the reconnect loop. @@ -325,10 +352,10 @@ func pumpConn(parentCtx context.Context, conn *transport.Conn, registry *agent.R return err } router, err := dispatch.New(dispatch.Config{ - Registry: registry, - Sender: conn, - Log: obslog.Bg(), - LocalWorkspace: local, + Registry: registry, + Sender: conn, + Log: obslog.Bg(), + Environments: localEnvironments(local), }) if err != nil { return fmt.Errorf("router init: %w", err) @@ -339,16 +366,11 @@ func pumpConn(parentCtx context.Context, conn *transport.Conn, registry *agent.R }() conn.StartHeartbeats(parentCtx, boot.HeartbeatInterval(), func() proto.HeartbeatPayload { - kinds := registry.SupportedAgentKinds() - for i := range kinds { - caps := &kinds[i].Capabilities - caps.WorkspaceOutputExport = proto.CapabilityFromBool(local.CanExport() && caps.LocalEnvironment.IsSupported() && caps.WorkspaceReadPreparation.IsSupported()) - } return proto.HeartbeatPayload{ Timestamp: time.Now().Unix(), ActiveRequests: router.ActiveRuns(), DaemonVersion: Version, - SupportedAgentKinds: kinds, + SupportedAgentKinds: localEnvironmentKinds(registry, local), HomeRemoval: proto.CapabilityUnsupported, } }, obslog.Bg().With("component", "heartbeat")) diff --git a/apps/daemon/internal/cli/connect_suspend.go b/apps/daemon/internal/cli/connect_suspend.go index 2fc97dba7..c449ee7c9 100644 --- a/apps/daemon/internal/cli/connect_suspend.go +++ b/apps/daemon/internal/cli/connect_suspend.go @@ -44,7 +44,7 @@ func newSuspendedRouter(conn *transport.Conn, registry *agent.Registry) (*suspen return nil, err } sender := &reconnectSender{conn: conn} - router, err := dispatch.New(dispatch.Config{Registry: registry, Sender: sender, Log: obslog.Bg(), LocalWorkspace: local}) + router, err := dispatch.New(dispatch.Config{Registry: registry, Sender: sender, Log: obslog.Bg(), Environments: localEnvironments(local)}) if err != nil { return nil, err } @@ -152,12 +152,7 @@ func (s *suspendedRouter) reconnectSuspension(ctx context.Context, dial transpor func (s *suspendedRouter) heartbeats(ctx context.Context, conn *transport.Conn, boot *transport.BootstrapResponse, discovery agentCLIDiscovery) { conn.StartHeartbeats(ctx, boot.HeartbeatInterval(), func() proto.HeartbeatPayload { - kinds := s.registry.SupportedAgentKinds() - for i := range kinds { - caps := &kinds[i].Capabilities - caps.WorkspaceOutputExport = proto.CapabilityFromBool(s.local.CanExport() && caps.LocalEnvironment.IsSupported() && caps.WorkspaceReadPreparation.IsSupported()) - } - return proto.HeartbeatPayload{Timestamp: time.Now().Unix(), ActiveRequests: s.router.ActiveRuns(), DaemonVersion: Version, SupportedAgentKinds: kinds, HomeRemoval: proto.CapabilityUnsupported} + return proto.HeartbeatPayload{Timestamp: time.Now().Unix(), ActiveRequests: s.router.ActiveRuns(), DaemonVersion: Version, SupportedAgentKinds: localEnvironmentKinds(s.registry, s.local), HomeRemoval: proto.CapabilityUnsupported} }, obslog.Bg()) } diff --git a/apps/daemon/internal/dispatch/assignment.go b/apps/daemon/internal/dispatch/assignment.go index de307618b..63e0e16f1 100644 --- a/apps/daemon/internal/dispatch/assignment.go +++ b/apps/daemon/internal/dispatch/assignment.go @@ -19,7 +19,9 @@ type assignmentState struct { // resource's Kind is empty when the bind carried none. resource sandboxbootstrap.Resource grant []byte - released bool + // environment is the owner resolved from the bind, or nil. + environment Environment + released bool // work counts the Session's admitted reads, writes, exports and Runtime // preparations until each has sent its terminal result. A release waits // for it, and the released assignment admits no more. @@ -49,6 +51,22 @@ func (r *Router) trackWorkLocked(ref proto.AssignmentRef) func() { return work.Done } +// admittedEnvironment admits ref for the Session and returns the owner that +// its assignment resolved, or nil, or the assignment rejection code. A +// Session's assignment, and so its owner, never changes on a Router. +func (r *Router) admittedEnvironment(ref proto.AssignmentRef, sessionID string) (Environment, string) { + r.mu.Lock() + defer r.mu.Unlock() + a := r.assignments[ref.SessionID] + if a == nil { + return nil, proto.AssignmentConflict + } + if code := r.admitLocked(ref, sessionID, a.environmentID); code != "" { + return nil, code + } + return a.environment, "" +} + // admitRunLocked admits a frame for the run, which ref must have started. // Router.mu must be held. func (r *Router) admitRunLocked(ref proto.AssignmentRef, state *sessionState) string { @@ -72,7 +90,14 @@ func (r *Router) handleAssignmentBind(ctx context.Context, env proto.Envelope) e a := r.assignments[ref.SessionID] switch { case a == nil: - r.assignments[ref.SessionID] = &assignmentState{ref: ref, environmentID: input.EnvironmentID, resource: resource, grant: input.AttachGrant} + var environment Environment + if r.environments != nil { + if environment = r.environments(ref, input); environment == nil { + code = proto.AssignmentConflict + break + } + } + r.assignments[ref.SessionID] = &assignmentState{ref: ref, environmentID: input.EnvironmentID, resource: resource, grant: input.AttachGrant, environment: environment} case a.ref.AssignmentID == ref.AssignmentID && (ref.Epoch < a.ref.Epoch || ref.Epoch == a.ref.Epoch && a.released): code = proto.AssignmentStale case a.ref != ref || a.environmentID != input.EnvironmentID || a.resource != resource || !bytes.Equal(a.grant, input.AttachGrant): @@ -84,8 +109,8 @@ func (r *Router) handleAssignmentBind(ctx context.Context, env proto.Envelope) e } // handleAssignmentRelease fences the assignment, then settles the Session's -// work and Executor and removes its home before it replies. A retry at the -// same epoch repeats the cleanup. +// work and Executor, closes its Environment owner and removes its home before +// it replies. A retry at the same epoch repeats the cleanup. func (r *Router) handleAssignmentRelease(ctx context.Context, env proto.Envelope) error { var input proto.AssignmentReleasePayload ref, code := env.Assignment, "" @@ -114,17 +139,22 @@ func (r *Router) handleAssignmentRelease(ctx context.Context, env proto.Envelope return r.reply(ctx, env, proto.TypeAssignmentStatus, assignmentStatus("", code)) } preparations := r.fenceSessionWorkLocked(ref.SessionID) - work := &r.assignments[ref.SessionID].work + work, environment := &r.assignments[ref.SessionID].work, r.assignments[ref.SessionID].environment r.shutdownWG.Add(1) r.mu.Unlock() go func() { defer r.shutdownWG.Done() + cleanupCtx, stop := r.shutdownContext(context.WithoutCancel(ctx)) + defer stop() for _, p := range preparations { r.releasePreparation(p, "failed", proto.AssignmentStale, true) } work.Wait() state, code := proto.AssignmentReleased, "" err := r.closeSessionExecutor(ref.SessionID) + if err == nil && environment != nil { + err = environment.Close(cleanupCtx) + } if err == nil && input.RemoveHome { state, err = proto.AssignmentHomeRemoved, r.removeHome(ref.SessionID) } @@ -132,9 +162,7 @@ func (r *Router) handleAssignmentRelease(ctx context.Context, env proto.Envelope r.log.Warn("assignment release cleanup unconfirmed", "session_id", ref.SessionID, "err", err) code = proto.CleanupUnconfirmed } - sendCtx, stop := r.shutdownContext(context.WithoutCancel(ctx)) - defer stop() - _ = r.reply(sendCtx, env, proto.TypeAssignmentStatus, assignmentStatus(state, code)) + _ = r.reply(cleanupCtx, env, proto.TypeAssignmentStatus, assignmentStatus(state, code)) }() return nil } diff --git a/apps/daemon/internal/dispatch/environment.go b/apps/daemon/internal/dispatch/environment.go index 3f8991935..2a845e0e9 100644 --- a/apps/daemon/internal/dispatch/environment.go +++ b/apps/daemon/internal/dispatch/environment.go @@ -1,11 +1,37 @@ package dispatch import ( + "context" "errors" + "io" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) +// Environment owns one Session's Environment: its resources and every effect on +// them. The Router resolves it from the Session's assignment_bind and admits, +// frames and fences each operation; the owner performs it. A Session without +// an owner declares none of these operations, and the Router rejects each with +// its typed unsupported code. docs/runtime-protocol.md defines the semantics. +type Environment interface { + // Configure checks an execution configuration against the Environment and + // returns it with the Environment's workspace root. It has no effects. + Configure(proto.PromptRequestPayload) (proto.PromptRequestPayload, error) + // Prepare fills the configured execution's installed capabilities before + // the Executor factory runs. + Prepare(context.Context, proto.PromptRequestPayload) (proto.PromptRequestPayload, error) + // ApplyRuntimePreparation applies one complete runtime_prepare transfer and + // returns only after its mutations stop. + ApplyRuntimePreparation(context.Context, proto.RuntimePreparePayload, []byte) error + ListWorkspaceDirectory(ctx context.Context, path string, maxEntries int) (agent.WorkspaceDirectoryResult, error) + WriteWorkspaceFile(ctx context.Context, path string, data []byte) (agent.WorkspaceWriteResult, error) + ExportOutputs(context.Context, io.Writer) error + // Close releases what the owner holds for the assignment once its work and + // Executors have settled. A retried release calls it again. + Close(context.Context) error +} + func validateExecutionEnvironment(req proto.PromptRequestPayload, caps proto.AgentKindCapabilities) error { if (req.LocalEnvironment != nil) == req.DisableExecutionEnvironment { return errors.New("execution requires exactly one of local_environment and disable_execution_environment") diff --git a/apps/daemon/internal/dispatch/environment_test.go b/apps/daemon/internal/dispatch/environment_test.go index e8d87d7a6..0a776ed5f 100644 --- a/apps/daemon/internal/dispatch/environment_test.go +++ b/apps/daemon/internal/dispatch/environment_test.go @@ -2,6 +2,8 @@ package dispatch_test import ( "context" + "crypto/sha256" + "encoding/hex" "errors" "sync/atomic" "testing" @@ -9,6 +11,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto/prototest" + "github.com/google/uuid" ) // assertPreparationOutcome waits for the terminal admission status of id. An @@ -92,3 +95,68 @@ func TestLocalEnvironmentRequiresAvailableCapability(t *testing.T) { }) } } + +// A Session whose assignment resolves no Environment owner declares no +// Environment operation: each gets its typed rejection before any effect. +func TestSessionWithoutOwnerRejectsEnvironmentOperations(t *testing.T) { + h := newHarness(t) + defer h.router.Shutdown(context.Background()) + var called atomic.Bool + registerExecutorKind(h.reg, proto.SupportedAgentKind{Kind: "local", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{LocalEnvironment: proto.CapabilitySupported})}, func(context.Context, proto.PromptRequestPayload) (agent.Executor, error) { + called.Store(true) + return nil, errors.New("controlled factory stop") + }) + assign(t, h.router, preparationSessionID, preparationEnvironmentID) + execution := proto.PromptRequestPayload{AgentKind: "local", AgentStateKey: stateKey(preparationSessionID), LocalEnvironment: &proto.LocalEnvironment{ID: preparationEnvironmentID}} + read := execution + read.WorkspaceReadOnly = true + for id, test := range map[string]struct { + request proto.PromptRequestPayload + code string + }{"read": {read, "unsupported_read_preparation"}, "execution": {execution, "invalid_configuration"}} { + _ = h.router.Handle(t.Context(), mustEnv(t, proto.TypeExecutionPrepare, id, proto.ExecutionPreparePayload{SessionID: preparationSessionID, Configuration: test.request})) + if status := waitPreparationStatus(t, h.sender, id, "rejected", ""); status.ErrorCode != test.code { + t.Fatalf("%s preparation = %+v", id, status) + } + } + empty := sha256.Sum256(nil) + for _, test := range []struct { + request, result string + payload any + code string + }{ + {proto.TypeRuntimePrepare, proto.TypeRuntimePrepareResult, proto.RuntimePreparePayload{Step: "begin", Action: "file", EnvironmentID: preparationEnvironmentID, SessionID: preparationSessionID, File: &proto.RuntimeInitialFile{Path: "/workspace/input"}, SHA256: hex.EncodeToString(empty[:])}, "runtime_preparation_unsupported"}, + {proto.TypeWorkspaceWrite, proto.TypeWorkspaceWriteResult, proto.WorkspaceWritePayload{Step: "begin", EnvironmentID: preparationEnvironmentID, SessionID: preparationSessionID, Path: "file", SHA256: hex.EncodeToString(empty[:])}, "write_unsupported"}, + {proto.TypeWorkspaceRead, proto.TypeWorkspaceReadResult, proto.WorkspaceReadPayload{Handle: "handle", EnvironmentID: preparationEnvironmentID, MaxEntries: 1}, "read_unsupported"}, + {proto.TypeWorkspaceExport, proto.TypeWorkspaceExportResult, proto.WorkspaceExportPayload{Step: "begin", Handle: "handle", EnvironmentID: preparationEnvironmentID}, "read_unsupported"}, + } { + id := uuid.NewString() + if err := h.router.Handle(t.Context(), mustEnv(t, test.request, id, test.payload)); err != nil { + t.Fatal(test.request, err) + } + waitFor(t, func() bool { return hasFrame(h.sender, test.result, id) }, test.result) + frame, _ := frameFor(h.sender, test.result, id) + var result struct { + Outcome string `json:"outcome"` + ErrorCode string `json:"error_code"` + } + if err := frame.DecodePayload(&result); err != nil || result.Outcome != "rejected" || result.ErrorCode != test.code { + t.Fatalf("%s = %+v, %v", test.request, result, err) + } + } + if called.Load() { + t.Fatal("an operation reached the Executor factory") + } +} + +func TestOwnedRuntimeRejectsForeignSessionBind(t *testing.T) { + h := localPreparationHarness(t) + defer h.router.Shutdown(context.Background()) + const foreign = "33333333-3333-4333-8333-333333333333" + if err := h.router.Handle(t.Context(), scoped(t, foreign, proto.TypeAssignmentBind, "bind-foreign", proto.AssignmentBindPayload{EnvironmentID: preparationEnvironmentID})); err != nil { + t.Fatal(err) + } + if status := waitAssignmentStatus(t, h.sender, "bind-foreign"); status.ErrorCode != proto.AssignmentConflict { + t.Fatalf("foreign bind = %+v", status) + } +} diff --git a/apps/daemon/internal/dispatch/executor.go b/apps/daemon/internal/dispatch/executor.go index 2fd2a18f5..ca4f5087d 100644 --- a/apps/daemon/internal/dispatch/executor.go +++ b/apps/daemon/internal/dispatch/executor.go @@ -57,10 +57,17 @@ func (r *Router) handleExecutorPrepare(ctx context.Context, env proto.Envelope, if err != nil || !caps.Preparation.IsSupported() { return r.rejectPreparation(env, "unsupported_preparation") } - if !r.sessionEnvironments { - if req, err = r.localWorkspace.Configure(req); err != nil { - return r.rejectPreparation(env, "invalid_configuration") - } + environment, code := r.admittedEnvironment(env.Assignment, input.SessionID) + if code != "" { + return r.rejectPreparation(env, code) + } + if environment != nil { + req, err = environment.Configure(req) + } else if req.LocalEnvironment != nil && !r.sessionEnvironments { + err = errors.New("the Session has no Environment owner") + } + if err != nil { + return r.rejectPreparation(env, "invalid_configuration") } if validateExecutionEnvironment(req, caps) != nil || len(req.FunctionTools) > 0 && !caps.FunctionTools.IsSupported() { return r.rejectPreparation(env, "unsupported_configuration") @@ -170,12 +177,12 @@ func (r *Router) handleExecutorPrepare(ctx context.Context, env proto.Envelope, if reused { r.publishPreparation(p, status) } else { - go r.prepareExecutor(p, req, factory) + go r.prepareExecutor(p, req, factory, environment) } return nil } -func (r *Router) prepareExecutor(p *preparationState, req proto.PromptRequestPayload, factory agent.ExecutorFactory) { +func (r *Router) prepareExecutor(p *preparationState, req proto.PromptRequestPayload, factory agent.ExecutorFactory, environment Environment) { defer r.shutdownWG.Done() owner := p.executor started := time.Now() @@ -188,8 +195,8 @@ func (r *Router) prepareExecutor(p *preparationState, req proto.PromptRequestPay var native agent.Executor var err error if owner.ctx.Err() == nil { - if !r.sessionEnvironments { - req, err = r.localWorkspace.Prepare(owner.ctx, req) + if environment != nil { + req, err = environment.Prepare(owner.ctx, req) } if err == nil && owner.ctx.Err() == nil { native, err = factory(owner.ctx, req) diff --git a/apps/daemon/internal/dispatch/export_test.go b/apps/daemon/internal/dispatch/export_test.go index 6ee90d36e..5a77ad2b1 100644 --- a/apps/daemon/internal/dispatch/export_test.go +++ b/apps/daemon/internal/dispatch/export_test.go @@ -1,5 +1,21 @@ package dispatch +import ( + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/localworkspace" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" +) + +// LocalEnvironments resolves binding as the owner of the Session it binds, as +// the daemon composes it. +func LocalEnvironments(binding *localworkspace.Binding) func(proto.AssignmentRef, proto.AssignmentBindPayload) Environment { + return func(ref proto.AssignmentRef, bind proto.AssignmentBindPayload) Environment { + if !binding.Matches(bind.EnvironmentID, ref.SessionID) { + return nil + } + return binding + } +} + func (r *Router) PreparationOwnershipForTest(handle string) (bool, bool) { r.mu.Lock() defer r.mu.Unlock() diff --git a/apps/daemon/internal/dispatch/local_directory_test.go b/apps/daemon/internal/dispatch/local_directory_test.go index d042e7854..019df9ba3 100644 --- a/apps/daemon/internal/dispatch/local_directory_test.go +++ b/apps/daemon/internal/dispatch/local_directory_test.go @@ -28,13 +28,13 @@ func TestLocalDirectoryPreparationNeedsNoHarnessAndRejectsOtherOwners(t *testing } var harnessCalls atomic.Int32 reg := agent.NewRegistry() - reg.RegisterKind(proto.SupportedAgentKind{Kind: "native", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{LocalEnvironment: proto.CapabilitySupported, WorkspaceReadPreparation: proto.CapabilitySupported})}, harnessconfig.Configuration{}) + reg.RegisterKind(proto.SupportedAgentKind{Kind: "native", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{LocalEnvironment: proto.CapabilitySupported})}, harnessconfig.Configuration{}) reg.RegisterExecutor("native", func(context.Context, proto.PromptRequestPayload) (agent.Executor, error) { harnessCalls.Add(1) return nil, errors.New("must not prepare a harness") }) sender := &recSender{} - r, err := dispatch.New(dispatch.Config{Registry: reg, Sender: sender, LocalWorkspace: binding}) + r, err := dispatch.New(dispatch.Config{Registry: reg, Sender: sender, Environments: dispatch.LocalEnvironments(binding)}) if err != nil { t.Fatal(err) } @@ -100,12 +100,12 @@ func TestLocalDirectoryKeepsNotDirectorySeparateFromFailures(t *testing.T) { t.Fatal(err) } reg := agent.NewRegistry() - reg.RegisterKind(proto.SupportedAgentKind{Kind: "native", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{LocalEnvironment: proto.CapabilitySupported, WorkspaceReadPreparation: proto.CapabilitySupported})}, harnessconfig.Configuration{}) + reg.RegisterKind(proto.SupportedAgentKind{Kind: "native", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{LocalEnvironment: proto.CapabilitySupported})}, harnessconfig.Configuration{}) reg.RegisterExecutor("native", func(context.Context, proto.PromptRequestPayload) (agent.Executor, error) { return nil, errors.New("must not prepare a harness") }) sender := &recSender{} - r, err := dispatch.New(dispatch.Config{Registry: reg, Sender: sender, LocalWorkspace: binding}) + r, err := dispatch.New(dispatch.Config{Registry: reg, Sender: sender, Environments: dispatch.LocalEnvironments(binding)}) if err != nil { t.Fatal(err) } diff --git a/apps/daemon/internal/dispatch/preparation.go b/apps/daemon/internal/dispatch/preparation.go index 81e51fabc..60b15d8ea 100644 --- a/apps/daemon/internal/dispatch/preparation.go +++ b/apps/daemon/internal/dispatch/preparation.go @@ -53,10 +53,16 @@ func (r *Router) handleExecutionPrepare(ctx context.Context, env proto.Envelope) if !caps.Preparation.IsSupported() { return r.rejectPreparation(env, "unsupported_preparation") } - if !caps.WorkspaceReadPreparation.IsSupported() || !proto.ValidWorkspaceReadPreparation(req) { + // The Session's owner declares read-only preparation for a kind that + // supports a local Environment. + environment, code := r.admittedEnvironment(env.Assignment, input.SessionID) + if code != "" { + return r.rejectPreparation(env, code) + } + if environment == nil || !caps.LocalEnvironment.IsSupported() || !proto.ValidWorkspaceReadPreparation(req) { return r.rejectPreparation(env, "unsupported_read_preparation") } - req, err := r.localWorkspace.Configure(req) + req, err := environment.Configure(req) if err != nil { return r.rejectPreparation(env, "invalid_configuration") } diff --git a/apps/daemon/internal/dispatch/preparation_test.go b/apps/daemon/internal/dispatch/preparation_test.go index bb39c151f..54109456c 100644 --- a/apps/daemon/internal/dispatch/preparation_test.go +++ b/apps/daemon/internal/dispatch/preparation_test.go @@ -103,7 +103,7 @@ func localPreparationHarness(t *testing.T) *harness { t.Fatal(err) } var err error - h.router, err = dispatch.New(dispatch.Config{Registry: h.reg, Sender: h.sender, LocalWorkspace: preparationWorkspace(t)}) + h.router, err = dispatch.New(dispatch.Config{Registry: h.reg, Sender: h.sender, Environments: dispatch.LocalEnvironments(preparationWorkspace(t))}) if err != nil { t.Fatal(err) } @@ -121,9 +121,9 @@ func preparationRequest() proto.ExecutionPreparePayload { func preparationRouter(t *testing.T, sender dispatch.Sender, timeout time.Duration, factory preparationFactory) *dispatch.Router { t.Helper() reg := agent.NewRegistry() - reg.RegisterKind(proto.SupportedAgentKind{Kind: "prepared", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{LocalEnvironment: proto.CapabilitySupported, WorkspaceReadPreparation: proto.CapabilitySupported, FunctionTools: proto.CapabilitySupported, Steering: proto.CapabilitySupported, DurableInputReceipts: proto.CapabilitySupported})}, prototest.ModelConfiguration()) + reg.RegisterKind(proto.SupportedAgentKind{Kind: "prepared", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{LocalEnvironment: proto.CapabilitySupported, FunctionTools: proto.CapabilitySupported, Steering: proto.CapabilitySupported, DurableInputReceipts: proto.CapabilitySupported})}, prototest.ModelConfiguration()) reg.RegisterExecutor("prepared", preparationExecutorFixture(factory)) - r, err := dispatch.New(dispatch.Config{Registry: reg, Sender: sender, PreparationTimeout: timeout, LocalWorkspace: preparationWorkspace(t)}) + r, err := dispatch.New(dispatch.Config{Registry: reg, Sender: sender, PreparationTimeout: timeout, Environments: dispatch.LocalEnvironments(preparationWorkspace(t))}) if err != nil { t.Fatal(err) } diff --git a/apps/daemon/internal/dispatch/router.go b/apps/daemon/internal/dispatch/router.go index 39cbf3065..3e0f68066 100644 --- a/apps/daemon/internal/dispatch/router.go +++ b/apps/daemon/internal/dispatch/router.go @@ -16,7 +16,6 @@ import ( "time" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/localworkspace" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" ) @@ -53,7 +52,7 @@ type Router struct { workspaceWrite *workspaceUpload workspaceExport *workspaceExport workspaceReads map[string]struct{} - localWorkspace *localworkspace.Binding + environments func(proto.AssignmentRef, proto.AssignmentBindPayload) Environment sessionEnvironments bool removeHome func(sessionID string) error } @@ -92,10 +91,14 @@ type Config struct { Log *slog.Logger IdleTimeout time.Duration PreparationTimeout time.Duration - LocalWorkspace *localworkspace.Binding - // SessionEnvironments says that a prepared execution's LocalEnvironment - // is its Session's Environment, which the Executor factory binds, and not - // a local workspace of this daemon. It excludes LocalWorkspace. + // Environments resolves the Environment owner of a Session's first bind on + // this Router, under the Router's lock. A nil owner rejects the bind: the + // Runtime does not serve that Session. Nil Environments leaves every + // Session without an owner. + Environments func(proto.AssignmentRef, proto.AssignmentBindPayload) Environment + // SessionEnvironments says that the Executor factory binds a prepared + // execution's LocalEnvironment itself, without an owner. It excludes + // Environments. SessionEnvironments bool // RemoveHome removes the Session's native home once its Executors have // closed. Nil declares that assignment_release does not accept RemoveHome. @@ -112,8 +115,8 @@ func New(cfg Config) (*Router, error) { if cfg.Sender == nil { return nil, errors.New("dispatch.New: Sender is required") } - if cfg.SessionEnvironments && cfg.LocalWorkspace != nil { - return nil, errors.New("dispatch.New: SessionEnvironments excludes LocalWorkspace") + if cfg.SessionEnvironments && cfg.Environments != nil { + return nil, errors.New("dispatch.New: SessionEnvironments excludes Environments") } log := cfg.Log if log == nil { @@ -141,7 +144,7 @@ func New(cfg Config) (*Router, error) { preparations: make(map[string]*preparationState), preparationRequests: make(map[string]*preparationState), preparationTimeout: preparationTimeout, - localWorkspace: cfg.LocalWorkspace, + environments: cfg.Environments, sessionEnvironments: cfg.SessionEnvironments, removeHome: cfg.RemoveHome, }, nil diff --git a/apps/daemon/internal/dispatch/runtime_preparation.go b/apps/daemon/internal/dispatch/runtime_preparation.go index 866eb710f..492282d4b 100644 --- a/apps/daemon/internal/dispatch/runtime_preparation.go +++ b/apps/daemon/internal/dispatch/runtime_preparation.go @@ -65,7 +65,12 @@ func (r *Router) handleRuntimePrepare(ctx context.Context, env proto.Envelope) e r.mu.Unlock() return r.sendRuntimePrepareResult(ctx, env, rejectedRuntimePreparation(code)) } - if r.localWorkspace == nil || !r.localWorkspace.Matches(request.EnvironmentID, request.SessionID) || r.runtimePreparationResourcesBusyLocked() { + environment := r.assignments[request.SessionID].environment + if environment == nil { + r.mu.Unlock() + return r.sendRuntimePrepareResult(ctx, env, rejectedRuntimePreparation("runtime_preparation_unsupported")) + } + if r.runtimePreparationResourcesBusyLocked() { r.mu.Unlock() return r.sendRuntimePrepareResult(ctx, env, rejectedRuntimePreparation("resource_unavailable")) } @@ -78,7 +83,7 @@ func (r *Router) handleRuntimePrepare(ctx context.Context, env proto.Envelope) e done := r.trackWorkLocked(env.Assignment) r.shutdownWG.Add(1) r.mu.Unlock() - go r.runRuntimePreparationTransfer(owner, u, r.localWorkspace.ApplyRuntimePreparation, done) + go r.runRuntimePreparationTransfer(owner, u, environment.ApplyRuntimePreparation, done) if err := r.sendRuntimePrepareResult(ctx, env, proto.RuntimePrepareResultPayload{Outcome: "ready"}); err != nil { cancel() return err diff --git a/apps/daemon/internal/dispatch/runtime_preparation_test.go b/apps/daemon/internal/dispatch/runtime_preparation_test.go index ca4497430..3a17d03cb 100644 --- a/apps/daemon/internal/dispatch/runtime_preparation_test.go +++ b/apps/daemon/internal/dispatch/runtime_preparation_test.go @@ -41,7 +41,7 @@ func capabilitiesTestRouter(t *testing.T) (*Router, *capabilitiesTestSender, str t.Fatal(err) } sender := &capabilitiesTestSender{frames: make(chan proto.Envelope, 64)} - router, err := New(Config{Registry: agent.NewRegistry(), Sender: sender, LocalWorkspace: binding}) + router, err := New(Config{Registry: agent.NewRegistry(), Sender: sender, Environments: LocalEnvironments(binding)}) if err != nil { t.Fatal(err) } @@ -150,7 +150,7 @@ func stringsOfZeroDigest() string { return hex.EncodeToString(make([]byte, sha25 func TestRuntimePreparationBeginRequiresExactBindingAndBounds(t *testing.T) { r, sender, environment, session := capabilitiesTestRouter(t) defer shutdownCapabilitiesRouter(t, r) - for _, mode := range []string{"environment", "session", "size", "nil-binding"} { + for _, mode := range []string{"environment", "session", "size", "no-owner"} { request := capabilityBegin(environment, session, []byte("abc")) switch mode { case "environment": @@ -159,8 +159,10 @@ func TestRuntimePreparationBeginRequiresExactBindingAndBounds(t *testing.T) { request.SessionID = uuid.NewString() case "size": request.SizeBytes = proto.RuntimePrepareMaxBytes + 1 - case "nil-binding": - r.localWorkspace = nil + case "no-owner": + r.mu.Lock() + r.assignments[session].environment = nil + r.mu.Unlock() } id := uuid.NewString() if err := r.Handle(t.Context(), capabilityEnvelope(t, id, request)); err != nil { diff --git a/apps/daemon/internal/dispatch/suspend_test.go b/apps/daemon/internal/dispatch/suspend_test.go index db35925ff..80e957c30 100644 --- a/apps/daemon/internal/dispatch/suspend_test.go +++ b/apps/daemon/internal/dispatch/suspend_test.go @@ -20,7 +20,11 @@ var suspendRef = proto.AssignmentRef{SessionID: "session", AssignmentID: "assign // bindAssignment records ref as bound in environmentID, as assignment_bind does. func bindAssignment(r *Router, ref proto.AssignmentRef, environmentID string) { r.mu.Lock() - r.assignments[ref.SessionID] = &assignmentState{ref: ref, environmentID: environmentID} + a := &assignmentState{ref: ref, environmentID: environmentID} + if r.environments != nil { + a.environment = r.environments(ref, proto.AssignmentBindPayload{EnvironmentID: environmentID}) + } + r.assignments[ref.SessionID] = a r.mu.Unlock() } diff --git a/apps/daemon/internal/dispatch/workspace_export.go b/apps/daemon/internal/dispatch/workspace_export.go index 724b84841..d460bef7b 100644 --- a/apps/daemon/internal/dispatch/workspace_export.go +++ b/apps/daemon/internal/dispatch/workspace_export.go @@ -47,9 +47,9 @@ func (r *Router) handleWorkspaceExport(ctx context.Context, env proto.Envelope) return errors.New("dispatch: workspace export request already pending") } } - _, code := r.workspaceResourceLocked(env.Assignment, proto.WorkspaceReadPayload{Handle: request.Handle, EnvironmentID: request.EnvironmentID}) + environment, code := r.workspaceResourceLocked(env.Assignment, proto.WorkspaceReadPayload{Handle: request.Handle, EnvironmentID: request.EnvironmentID}) p := r.preparations[request.Handle] - if code == "" && (u != nil || r.workspaceWrite != nil || !r.localWorkspace.CanExport() || p == nil || p.executor != nil) { + if code == "" && (u != nil || r.workspaceWrite != nil || p.executor != nil) { code = "resource_unavailable" } if code != "" { @@ -64,13 +64,13 @@ func (r *Router) handleWorkspaceExport(ctx context.Context, env proto.Envelope) done := r.trackWorkLocked(env.Assignment) r.shutdownWG.Add(1) r.mu.Unlock() - go r.runWorkspaceExport(owner, u, done) + go r.runWorkspaceExport(owner, u, environment, done) return nil } // runWorkspaceExport answers each request it admitted, even once the export is // canceled, and runs done after its last result is sent. -func (r *Router) runWorkspaceExport(ctx context.Context, u *workspaceExport, done func()) { +func (r *Router) runWorkspaceExport(ctx context.Context, u *workspaceExport, environment Environment, done func()) { defer r.shutdownWG.Done() defer done() // A result has its own send budget, independent of the export's cancellation. @@ -81,7 +81,7 @@ func (r *Router) runWorkspaceExport(ctx context.Context, u *workspaceExport, don exported := make(chan struct{}) go func() { defer close(exported) - err := r.localWorkspace.ExportOutputs(ctx, writer) + err := environment.ExportOutputs(ctx, writer) _ = writer.CloseWithError(err) }() var offset int64 diff --git a/apps/daemon/internal/dispatch/workspace_export_test.go b/apps/daemon/internal/dispatch/workspace_export_test.go index cb142173f..b398db252 100644 --- a/apps/daemon/internal/dispatch/workspace_export_test.go +++ b/apps/daemon/internal/dispatch/workspace_export_test.go @@ -56,7 +56,7 @@ func exporterRouter(t *testing.T, program string) (*Router, exportSender, proto. t.Fatal(err) } sender := exportSender{make(chan proto.Envelope, 8)} - r, err := New(Config{Registry: agent.NewRegistry(), Sender: sender, LocalWorkspace: binding}) + r, err := New(Config{Registry: agent.NewRegistry(), Sender: sender, Environments: LocalEnvironments(binding)}) if err != nil { t.Fatal(err) } diff --git a/apps/daemon/internal/dispatch/workspace_read.go b/apps/daemon/internal/dispatch/workspace_read.go index d2ef331c7..f9ff2f561 100644 --- a/apps/daemon/internal/dispatch/workspace_read.go +++ b/apps/daemon/internal/dispatch/workspace_read.go @@ -36,7 +36,7 @@ func (r *Router) handleWorkspaceRead(ctx context.Context, env proto.Envelope) er r.mu.Unlock() return r.sendWorkspaceRead(ctx, env, rejectedWorkspaceRead("read_capacity")) } - lister, code := r.workspaceResourceLocked(env.Assignment, request) + environment, code := r.workspaceResourceLocked(env.Assignment, request) if code != "" { r.mu.Unlock() return r.sendWorkspaceRead(ctx, env, rejectedWorkspaceRead(code)) @@ -55,45 +55,40 @@ func (r *Router) handleWorkspaceRead(ctx context.Context, env proto.Envelope) er // Observer loss does not discard an admitted native wait or replay it. operation, cancel := context.WithTimeout(context.WithoutCancel(ctx), 12*time.Second) defer cancel() - result := listWorkspaceDirectory(operation, lister, request) + result := listWorkspaceDirectory(operation, environment, request) _ = r.sendWorkspaceRead(context.WithoutCancel(ctx), env, result) }() return nil } -// workspaceResourceLocked returns what lists ref's directory: the local -// workspace, or without one the Harness Session of the Run ref admitted. -// Only the local workspace serves a preparation handle; a read-only -// preparation needs it. -func (r *Router) workspaceResourceLocked(ref proto.AssignmentRef, request proto.WorkspaceReadPayload) (agent.WorkspaceDirectoryLister, string) { +// workspaceResourceLocked returns the Environment owner of ref's Session once +// ref admits the read's ready preparation handle or open Run. +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 } + environment := r.assignments[ref.SessionID].environment + 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) || r.localWorkspace == nil { + !p.owns || p.busy || p.ctx.Err() != nil || !time.Now().Before(p.deadline) { return nil, "resource_unavailable" } - return r.localWorkspace, "" + return environment, "" } s := r.sessions[request.RunID] if s == nil || s.assignment != ref || s.environmentID != request.EnvironmentID || s.session == nil || !r.runRouteOpenLocked(s) { return nil, "resource_unavailable" } - if r.localWorkspace != nil { - return r.localWorkspace, "" - } - lister, ok := s.session.(agent.WorkspaceDirectoryLister) - if !ok { - return nil, "read_unsupported" - } - return lister, "" + return environment, "" } -func listWorkspaceDirectory(ctx context.Context, lister agent.WorkspaceDirectoryLister, request proto.WorkspaceReadPayload) proto.WorkspaceReadResultPayload { - read, err := lister.ListWorkspaceDirectory(ctx, request.Path, request.MaxEntries) +func listWorkspaceDirectory(ctx context.Context, environment Environment, request proto.WorkspaceReadPayload) proto.WorkspaceReadResultPayload { + read, err := environment.ListWorkspaceDirectory(ctx, request.Path, request.MaxEntries) if err != nil { return workspaceReadFailure(err) } @@ -119,9 +114,7 @@ func workspaceReadFailure(err error) proto.WorkspaceReadResultPayload { err error code string }{ - {agent.ErrWorkspaceReadUnsupported, "read_unsupported"}, {agent.ErrWorkspaceReadUnavailable, "resource_unavailable"}, - {agent.ErrWorkspaceReadBusy, "read_capacity"}, {agent.ErrWorkspaceReadInvalid, "invalid_request"}, {agent.ErrWorkspaceNotDirectory, proto.WorkspaceReadNotDirectory}, {fs.ErrNotExist, "not_found"}, diff --git a/apps/daemon/internal/dispatch/workspace_write.go b/apps/daemon/internal/dispatch/workspace_write.go index f3f80e968..e175cbe3c 100644 --- a/apps/daemon/internal/dispatch/workspace_write.go +++ b/apps/daemon/internal/dispatch/workspace_write.go @@ -12,7 +12,7 @@ import ( "github.com/google/uuid" ) -// Router.mu protects this single bounded transfer for the dedicated Environment. +// Router.mu protects this single bounded transfer to an Environment owner. type workspaceUpload struct { envelope proto.Envelope request proto.WorkspaceWritePayload @@ -62,7 +62,12 @@ func (r *Router) handleWorkspaceWrite(ctx context.Context, env proto.Envelope) e r.mu.Unlock() return r.sendWorkspaceWrite(ctx, env, rejectedWorkspaceWrite(code)) } - if !r.localWorkspace.AcceptsFileWrite(request.EnvironmentID, request.SessionID) || len(r.sessions) != 0 || len(r.workspaceReads) != 0 { + environment := r.assignments[request.SessionID].environment + if environment == nil { + r.mu.Unlock() + return r.sendWorkspaceWrite(ctx, env, rejectedWorkspaceWrite("write_unsupported")) + } + if len(r.sessions) != 0 || len(r.workspaceReads) != 0 { r.mu.Unlock() return r.sendWorkspaceWrite(ctx, env, rejectedWorkspaceWrite("resource_unavailable")) } @@ -83,7 +88,7 @@ func (r *Router) handleWorkspaceWrite(ctx context.Context, env proto.Envelope) e done := r.trackWorkLocked(env.Assignment) r.shutdownWG.Add(1) r.mu.Unlock() - go r.runWorkspaceUpload(context.WithoutCancel(ctx), u, done) + go r.runWorkspaceUpload(context.WithoutCancel(ctx), u, environment, done) return r.sendWorkspaceWrite(ctx, env, proto.WorkspaceWriteResultPayload{Outcome: "ready"}) } u := r.workspaceWrite @@ -111,7 +116,7 @@ func (r *Router) handleWorkspaceWrite(ctx context.Context, env proto.Envelope) e return nil } -func (r *Router) runWorkspaceUpload(ctx context.Context, u *workspaceUpload, done func()) { +func (r *Router) runWorkspaceUpload(ctx context.Context, u *workspaceUpload, environment Environment, done func()) { defer r.shutdownWG.Done() defer done() timer := time.NewTimer(120 * time.Second) @@ -134,7 +139,7 @@ func (r *Router) runWorkspaceUpload(ctx context.Context, u *workspaceUpload, don result = rejectedWorkspaceWrite(fenced) } if apply { - write, err := r.localWorkspace.WriteWorkspaceFile(ctx, u.request.Path, data) + write, err := environment.WriteWorkspaceFile(ctx, u.request.Path, data) result = workspaceWriteResult(write, err, u.request.SizeBytes) } r.mu.Lock() @@ -169,7 +174,6 @@ func workspaceWriteResult(write agent.WorkspaceWriteResult, err error, size int) err error code string }{ - {agent.ErrWorkspaceWriteUnsupported, "write_unsupported"}, {agent.ErrWorkspaceWriteUnavailable, "resource_unavailable"}, {agent.ErrWorkspaceWriteBusy, "write_capacity"}, {agent.ErrWorkspaceWriteInvalid, "invalid_request"}, diff --git a/apps/daemon/internal/dispatch/workspace_write_test.go b/apps/daemon/internal/dispatch/workspace_write_test.go index 7ee243ed0..f4964ad97 100644 --- a/apps/daemon/internal/dispatch/workspace_write_test.go +++ b/apps/daemon/internal/dispatch/workspace_write_test.go @@ -26,7 +26,7 @@ func localWriterRouter(t *testing.T) (*dispatch.Router, *recSender, proto.Worksp t.Fatal(err) } sender := &recSender{} - r, err := dispatch.New(dispatch.Config{Registry: agent.NewRegistry(), Sender: sender, LocalWorkspace: binding}) + r, err := dispatch.New(dispatch.Config{Registry: agent.NewRegistry(), Sender: sender, Environments: dispatch.LocalEnvironments(binding)}) if err != nil { t.Fatal(err) } diff --git a/apps/daemon/internal/localworkspace/binding.go b/apps/daemon/internal/localworkspace/binding.go index 0ead8710a..c3110e3d2 100644 --- a/apps/daemon/internal/localworkspace/binding.go +++ b/apps/daemon/internal/localworkspace/binding.go @@ -1,6 +1,7 @@ package localworkspace import ( + "context" "errors" "os" "strings" @@ -55,10 +56,7 @@ func Load() (*Binding, error) { // Configure validates the reference before supplying the immutable local cwd. func (b *Binding) Configure(r proto.PromptRequestPayload) (proto.PromptRequestPayload, error) { - if b == nil && r.LocalEnvironment == nil { - return r, nil - } - if b == nil || r.LocalEnvironment == nil || r.LocalEnvironment.ID != b.environment || r.AgentStateKey != b.stateKey || + if r.LocalEnvironment == nil || r.LocalEnvironment.ID != b.environment || r.AgentStateKey != b.stateKey || r.DisableExecutionEnvironment { return r, errors.New("request does not match the dedicated local Environment") } @@ -91,3 +89,6 @@ func (b *Binding) Configure(r proto.PromptRequestPayload) (proto.PromptRequestPa func (b *Binding) Matches(environment, session string) bool { return b != nil && b.environment == environment && b.stateKey == "agents-api-"+session } + +// Close keeps the workspace, which outlives each assignment of its Session. +func (b *Binding) Close(context.Context) error { return nil } diff --git a/apps/daemon/internal/localworkspace/binding_test.go b/apps/daemon/internal/localworkspace/binding_test.go index b23ea8c0f..68a30dad6 100644 --- a/apps/daemon/internal/localworkspace/binding_test.go +++ b/apps/daemon/internal/localworkspace/binding_test.go @@ -50,12 +50,6 @@ func TestBindingRejectsScopeOverrides(t *testing.T) { } }) } - if _, err := (*Binding)(nil).Configure(valid); err == nil { - t.Fatal("unbound Runtime accepted a local Environment") - } - if _, err := (*Binding)(nil).Configure(proto.PromptRequestPayload{}); err != nil { - t.Fatal("ordinary unbound behavior changed", err) - } } func TestDirectoryValidatesRelativePaths(t *testing.T) { diff --git a/apps/daemon/internal/localworkspace/export.go b/apps/daemon/internal/localworkspace/export.go index a853e7aaf..c4237105e 100644 --- a/apps/daemon/internal/localworkspace/export.go +++ b/apps/daemon/internal/localworkspace/export.go @@ -3,15 +3,12 @@ package localworkspace import ( "context" "errors" - "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "io" + + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) -func (b *Binding) CanExport() bool { return b != nil && b.workspace != "" } func (b *Binding) ExportOutputs(ctx context.Context, output io.Writer) error { - if !b.CanExport() || output == nil { - return errors.New("workspace export unavailable") - } return b.exportNativeOutputs(ctx, &exportWriter{output: output}) } diff --git a/apps/daemon/internal/localworkspace/write.go b/apps/daemon/internal/localworkspace/write.go index 05e7a6dc2..8e14d0654 100644 --- a/apps/daemon/internal/localworkspace/write.go +++ b/apps/daemon/internal/localworkspace/write.go @@ -13,18 +13,9 @@ import ( // This private transfer bound is distinct from the public inline-file limit. const WriteMaxBytes = proto.WorkspaceWriteMaxBytes -var _ agent.WorkspaceWriter = (*Binding)(nil) - -func (b *Binding) AcceptsFileWrite(environment, session string) bool { - return b != nil && b.writer != nil && b.environment == environment && b.stateKey == "agents-api-"+session -} - // WriteWorkspaceFile starts only after the caller supplies the complete bounded // body. Core must persist mutation ownership before invoking this operation. func (b *Binding) WriteWorkspaceFile(ctx context.Context, path string, data []byte) (result agent.WorkspaceWriteResult, err error) { - if b == nil || b.writer == nil { - return result, agent.ErrWorkspaceWriteUnsupported - } if len(data) > WriteMaxBytes || len(path) > 4096 || path == "." || !fs.ValidPath(path) || strings.ContainsAny(path, "\\\x00\r\n") { return result, agent.ErrWorkspaceWriteInvalid } diff --git a/contracts/agents-api/harness-onboarding.md b/contracts/agents-api/harness-onboarding.md index 661e4dd47..25522f1e4 100644 --- a/contracts/agents-api/harness-onboarding.md +++ b/contracts/agents-api/harness-onboarding.md @@ -58,7 +58,7 @@ Implement the mandatory text lifecycle and handle every extension explicitly. Qu ## Required adapter interfaces -[`agent/harness.go`](https://github.com/MiniMax-AI/OpenAgentCore/blob/main/apps/daemon/internal/agent/harness.go) is the interface entry point. The required lifecycle is `ExecutorFactory`, `Executor`, `Turn` (including `DurableSteerer`) and `TurnSettlement`. Required methods perform their native obligations; returning Unsupported is not an implementation of cancellation, receipts, settlement or cleanup. Turn and workspace extension interfaces stay small and separate, but every public adapter implements each one explicitly. All use the neutral protocol types. +[`agent/harness.go`](https://github.com/MiniMax-AI/OpenAgentCore/blob/main/apps/daemon/internal/agent/harness.go) is the interface entry point. The required lifecycle is `ExecutorFactory`, `Executor`, `Turn` (including `DurableSteerer`) and `TurnSettlement`. Required methods perform their native obligations; returning Unsupported is not an implementation of cancellation, receipts, settlement or cleanup. Turn extension interfaces stay small and separate, but every public adapter implements each one explicitly. All use the neutral protocol types. For example, the Codex adapter keeps its app-server and thread, the Claude adapter one streaming Query, and the MiniMax adapter its ACP connection and native session. All expose the same Executor and Turn contract. Native callbacks and resources stay inside the adapter; the Runtime owns admission, idle expiry and replacement. Cancellation targets the exact Turn through `agent.Session`, and the adapter supplies native completion evidence to the Runtime. @@ -69,7 +69,6 @@ For example, the Codex adapter keeps its app-server and thread, the Claude adapt | `DurableSteerer` | Real implementation on every Turn | Distinguish a complete write from the native application receipt; keep retry identity | | `Steerer` | Explicit implementation or Unsupported | Additional non-durable active-Turn input | | `FunctionResultSubmitter` | Explicit implementation or Unsupported | Match native call and result identity and acknowledge application | -| `WorkspaceDirectoryLister`, `WorkspaceWriter` | Explicit on Turn and Executor owners | Use the authorized workspace, confirm access, commit or close, or return the operation's Unsupported error | | Neutral messages, images, MCP, structured output and Subagent observations | Explicit capability decisions | Keep each operation's protocol semantics; reject unsupported input before submission | Each adapter's `contracts.go` holds an individual compile-time assertion for each small interface. Do not embed a default implementation that makes future interfaces appear implemented. Adding a contract also requires a classification in the common completeness check and an explicit assertion in every public adapter; the check follows the authored Harness catalog. @@ -86,7 +85,7 @@ The reason is a fixed safe string, never submitted content, a credential or raw The wire request carries no working directory. The Runtime checks `local_environment.workspace_directory` against its binding and gives the Harness its bound workspace directory in `LocalEnvironment.WorkspaceRoot`; run the native Harness there. -Workspace capability describes the actual Runtime and resource-owner combination. The Codex and MiniMax resource objects reject native workspace access while the common authorized `localworkspace` owner provides it; Claude can expose native read and list access, and the common owner provides writes. Interface presence alone never selects a resource or advertises support. +Workspace reads, writes, output export and read-only preparation belong to the Session's [Environment owner](../../docs/runtime-protocol.md#session-assignments), not the adapter. An adapter implements none of them and declares `WorkspaceReadPreparation` and `WorkspaceOutputExport` unsupported; the Runtime sets both from its owner. The service profile qualifies public combinations and the Runtime advertises the installed combination; neither replaces schema validation or Project authorization. Native behavior tests must agree with the declarations. An advertised operation that returns Unsupported is a contract violation, never success or grounds for replay. @@ -161,7 +160,7 @@ Registration is static and requires a build. Export one `agent.Declaration` from Every `proto.AgentKindCapabilities` field must be explicitly `proto.CapabilitySupported` or `proto.CapabilityUnsupported`, even for an unavailable Harness. `proto.CapabilityUnspecified` is invalid: zero values and omitted fields never mean Unsupported. An installation probe may set an individual field with `proto.CapabilityFromBool`; it must not populate unmentioned or future fields. Availability stays separate in `SupportedAgentKind.Available`. Registration validates the complete declaration before changing the registry, and the wire carries an explicit boolean for every field, so omitted and null fields are invalid. A new field requires a decision in every production declaration. Runtime consumers use `IsSupported()` and reject unsupported requests before native operations; an interface assertion verifies implementation, never support. Every declaration must match the behavior verified for that installation; the [Core–Runtime protocol](../../docs/runtime-protocol.md#capability-declarations) owns how declarations travel and are frozen. -The admission mapping is explicit. `Steering` controls non-durable `Steerer` input. `DurableInputReceipts` controls `DurableSteerer` input and also requires the Turn settlement contract; neither implies the other, and Core's public text profile requires both. Workspace declarations describe the authorized resource owner, including the common Runtime workspace implementation. `WorkspaceReadPreparation` admits `execution_prepare` with `workspace_read_only`. The Runtime readies that preparation itself and serves its reads from the bound local workspace directory without calling the Executor factory, so an adapter declares it only when this kind's Environment workspace is that local directory. Registration does not derive it. Runtime registration does not grant Core qualification; the service profile does. +The admission mapping is explicit. `Steering` controls non-durable `Steerer` input. `DurableInputReceipts` controls `DurableSteerer` input and also requires the Turn settlement contract; neither implies the other, and Core's public text profile requires both. `WorkspaceReadPreparation` admits `execution_prepare` with `workspace_read_only`, which the Environment owner readies and serves without calling the Executor factory. Runtime registration does not grant Core qualification; the service profile does. The runnable test-only example [`testdata/onboarding/main.go`](https://github.com/MiniMax-AI/OpenAgentCore/blob/main/apps/daemon/testdata/onboarding/main.go) registers a text-only synthetic Harness under the `mcode` kind, because Core admits only [catalog](./harness-catalog.md) Harnesses. It shows a Session-owned Executor, fresh Turns, durable steering, cancellation and history binding, and is never shipped. diff --git a/contracts/agents-api/zh/harness-onboarding.md b/contracts/agents-api/zh/harness-onboarding.md index 5d7b26f55..28b3fc293 100644 --- a/contracts/agents-api/zh/harness-onboarding.md +++ b/contracts/agents-api/zh/harness-onboarding.md @@ -1,7 +1,7 @@ --- title: "将原生 Harness 添加到 OpenAgentCore" source: contracts/agents-api/harness-onboarding.md -source_hash: 6c87f33d0ea60e0510c262b58aa92279852457bb85f0443b87a38bf2b0c3419c +source_hash: 146f98f62e43880478349005b7dcb8dd91d03f26579b38ef3f025f158b722a1c --- **Harness** 是一种运行模型和工具循环的原生代理引擎(Codex、Claude Code、MiniMax Code)。**Harness 适配器**将 Runtime 的 Executor 和 Turn 契约转换到该引擎的 SDK 或协议。本文档定义 Runtime–Harness 协议:适配器接口及其生命周期义务、注册、Core 资格认定和验收。[Harness capabilities](harness-capabilities.md) 记录了当前每个 Harness 支持的功能。 @@ -60,7 +60,7 @@ Environment 提供执行资源。受管 E2B、Docker 和 microsandbox 机器以 ## 必需的适配器接口 {#required-adapter-interfaces} -[`agent/harness.go`](https://github.com/MiniMax-AI/OpenAgentCore/blob/main/apps/daemon/internal/agent/harness.go) 是接口入口。必需的生命周期包括 `ExecutorFactory`、`Executor`、`Turn`(包括 `DurableSteerer`)和 `TurnSettlement`。必需方法必须履行其原生义务;返回 Unsupported 并不构成对取消、回执、结算或清理的实现。Turn 和工作区扩展接口应保持小而独立,但每个公共适配器都必须明确实现每一个接口。所有接口都使用中立协议类型。 +[`agent/harness.go`](https://github.com/MiniMax-AI/OpenAgentCore/blob/main/apps/daemon/internal/agent/harness.go) 是接口入口。必需的生命周期包括 `ExecutorFactory`、`Executor`、`Turn`(包括 `DurableSteerer`)和 `TurnSettlement`。必需方法必须履行其原生义务;返回 Unsupported 并不构成对取消、回执、结算或清理的实现。Turn 扩展接口应保持小而独立,但每个公共适配器都必须明确实现每一个接口。所有接口都使用中立协议类型。 例如,Codex 适配器保留其 app-server 和 thread,Claude 适配器保留一个流式 Query,MiniMax 适配器保留其 ACP 连接和原生 session。它们都公开相同的 Executor 和 Turn 契约。原生回调和资源保留在适配器内部;Runtime 负责准入、空闲过期和替换。取消通过 `agent.Session` 精确定位到目标 Turn,适配器则向 Runtime 提供原生完成证据。 @@ -71,7 +71,6 @@ Environment 提供执行资源。受管 E2B、Docker 和 microsandbox 机器以 | `DurableSteerer` | 每个 Turn 上真实实现 | 区分完整写入与原生应用回执;保留重试身份 | | `Steerer` | 明确实现或 Unsupported | 额外的非持久化活动 Turn 输入 | | `FunctionResultSubmitter` | 明确实现或 Unsupported | 匹配原生调用和结果身份,并确认应用 | -| `WorkspaceDirectoryLister`、`WorkspaceWriter` | 在 Turn 和 Executor 所有者上明确实现 | 使用授权工作区,确认访问,提交或关闭,或者返回该操作的 Unsupported 错误 | | 中立消息、图像、MCP、结构化输出和 Subagent 观察 | 明确作出能力决策 | 保持每项操作的协议语义;在提交前拒绝不受支持的输入 | 每个适配器的 `contracts.go` 都包含针对每个小型接口的单项编译时断言。不要嵌入会让未来接口看起来已经实现的默认实现。添加契约时,还必须在通用完整性检查中进行分类,并在每个公共适配器中添加明确断言;该检查遵循已编写的 Harness 目录。 @@ -88,7 +87,7 @@ func (s *Session) SubmitFunctionResult(context.Context, proto.FunctionResultPayl 线协议请求不携带工作目录。Runtime 将 `local_environment.workspace_directory` 与其绑定进行核对,并通过 `LocalEnvironment.WorkspaceRoot` 向 Harness 提供其绑定的工作区目录;必须在该目录中运行原生 Harness。 -工作区能力描述实际 Runtime 与资源所有者的组合。Codex 和 MiniMax 资源对象拒绝原生工作区访问,而通用的授权 `localworkspace` 所有者提供该访问;Claude 可以公开原生读取和列举访问,通用所有者提供写入。仅仅存在相应接口绝不会选择某个资源或宣称支持。 +工作区读取、写入、输出导出和只读 preparation 属于 Session 的 [Environment owner](../../../docs/zh/runtime-protocol.md#session-assignments),不属于 adapter。adapter 不实现其中任何操作,并将 `WorkspaceReadPreparation` 和 `WorkspaceOutputExport` 声明为不支持;Runtime 根据其 owner 设置二者。 服务 profile 对公共组合进行资格认定,Runtime 宣称已安装的组合;二者都不能替代 schema 验证或 Project 授权。原生行为测试必须与声明一致。已宣称但返回 Unsupported 的操作属于契约违规,既不是成功,也不能作为重放的依据。 @@ -163,7 +162,7 @@ MCP、公共函数、延迟函数发现、结构化输出、图像输入、详 每个 `proto.AgentKindCapabilities` 字段都必须显式设为 `proto.CapabilitySupported` 或 `proto.CapabilityUnsupported`,即使 Harness 不可用也是如此。`proto.CapabilityUnspecified` 无效:零值和省略字段绝不表示 Unsupported。安装探测可以使用 `proto.CapabilityFromBool` 设置单个字段;但不得填充未提及字段或未来字段。可用性通过 `SupportedAgentKind.Available` 单独表示。注册会在更改 registry 之前验证完整声明;线协议会为每个字段携带显式布尔值,因此省略字段和 null 字段均无效。添加新字段时,每个生产声明都必须作出决定。Runtime 使用者应调用 `IsSupported()`,并在原生操作前拒绝不受支持的请求;接口断言用于验证实现,绝不表示支持。每个声明都必须与针对该安装验证的行为一致;[Core–Runtime protocol](../../../docs/zh/runtime-protocol.md#capability-declarations) 负责声明的传输方式和冻结方式。 -准入映射是显式的。`Steering` 控制非持久化 `Steerer` 输入。`DurableInputReceipts` 控制 `DurableSteerer` 输入,并且还要求 Turn 结算契约;二者互不隐含,而且 Core 的公共文本 profile 要求同时具备二者。工作区声明描述授权资源所有者,包括通用 Runtime 工作区实现。`WorkspaceReadPreparation` 准入带 `workspace_read_only` 的 `execution_prepare`。Runtime 自行就绪该 preparation,并从绑定的本地工作区目录提供读取,不调用 Executor 工厂;因此仅当该 kind 的 Environment 工作区就是该本地目录时,adapter 才声明它。注册过程不会派生它。Runtime 注册不会授予 Core 资格;服务 profile 才会授予。 +准入映射是显式的。`Steering` 控制非持久化 `Steerer` 输入。`DurableInputReceipts` 控制 `DurableSteerer` 输入,并且还要求 Turn 结算契约;二者互不隐含,而且 Core 的公共文本 profile 要求同时具备二者。`WorkspaceReadPreparation` 准入带 `workspace_read_only` 的 `execution_prepare`,由 Environment owner 就绪并提供读取,不调用 Executor 工厂。Runtime 注册不会授予 Core 资格;服务 profile 才会授予。 可运行的仅测试示例 [`testdata/onboarding/main.go`](https://github.com/MiniMax-AI/OpenAgentCore/blob/main/apps/daemon/testdata/onboarding/main.go) 会以 `mcode` 类型注册一个仅支持文本的合成 Harness,因为 Core 只接纳[目录](harness-catalog.md)中的 Harness。它展示 Session 所有的 Executor、全新的 Turn、持久化引导、取消和历史绑定,并且绝不会发布。 diff --git a/docs/runtime-protocol.md b/docs/runtime-protocol.md index 3ccea3be1..48dbef51f 100644 --- a/docs/runtime-protocol.md +++ b/docs/runtime-protocol.md @@ -28,7 +28,7 @@ Each physical connection has fresh routing, admission handles and transfer state On the wire each field is a JSON boolean, and every field is present, including `false`. Encoding an incomplete declaration fails. Decoding rejects omitted, null, invalid and unknown fields, and a missing capability object. An invalid heartbeat clears the connection's admission snapshot and closes its transport; that establishes no native completion or cancellation result. -Each admitted Executor and Turn keeps the declaration it was admitted with. A later heartbeat cannot add operations to an existing owner. Optional operations check this snapshot before any native call; the presence of a Go interface never grants support. A declared operation that returns `agent.ErrUnsupportedOperation` is a contract violation, distinct from unavailability, a failed native call or an uncertain write. Uncertain operations keep their receipts and ownership and are never replayed automatically. Workspace support includes the common Runtime workspace implementation, so a native adapter's unsupported workspace method does not disable that composition. +Each admitted Executor and Turn keeps the declaration it was admitted with. A later heartbeat cannot add operations to an existing owner. Optional operations check this snapshot before any native call; the presence of a Go interface never grants support. A declared operation that returns `agent.ErrUnsupportedOperation` is a contract violation, distinct from unavailability, a failed native call or an uncertain write. Uncertain operations keep their receipts and ownership and are never replayed automatically. A Runtime declares `workspace_read_preparation` and `workspace_output_export` from its [Environment owner](#session-assignments), never from a Harness adapter. A new field requires an explicit decision in every production declaration. Contract tests enumerate every field for registration, wire round trips and the persisted boolean projection; the shared test fixture lists fields individually and supplies no defaults for future ones. [Harness onboarding](../contracts/agents-api/harness-onboarding.md) owns the adapter side of each declaration. @@ -116,7 +116,9 @@ Every Session frame carries the assignment: `execution_prepare`, `execution_star Before a Session's first operation on a connection, including Environment initialization and file work without a Turn, Core sends `assignment_bind` with the Session's Environment ID and waits for `assignment_status` `bound`. When the Runtime is an agent host and the Environment has a live [Link](./sandbox-link-protocol.md) resource, the bind also carries `resource`, that resource as the [bootstrap input](./sandbox-bootstrap.md#launch-input) names it, and `attach_grant`, the base64 grant with which the agent host opens services on that resource generation under this assignment and epoch. The grant is secret. Core sends neither field to any other Runtime. A bind with only one of them, or with a resource of another Environment, fails with `invalid_request`. A repeated bind of the same assignment with the same Environment, resource and grant is `bound` again; any other bind of it fails with `assignment_conflict`. The Runtime admits a Session frame only under the assignment it bound: an older epoch, or a released one, fails with `assignment_stale`; another assignment, Session or Environment fails with `assignment_conflict`. A started Run's frames, including its cancellation receipt, stay admissible under the assignment that started it until the release. A repeated function result or decision whose receipt the Runtime already recorded is answered only under the assignment that applied it; another fails with `assignment_conflict`. -Core records a release and advances the epoch before it sends anything, which withdraws the assignment's attach grant; it then has the relay revoke the Environment's Link resource at its current generation, so the attachments opened under the grant close before the Runtime receives the release. Deleting a Session releases its assignment with `remove_home: true`; releasing its Environment sends `false`. A deletion never revokes a shared Runtime credential. `assignment_release` fences the assignment at once. The Runtime then stops the Session's work: a transfer still receiving its body, or committed but not yet applied, ends with `assignment_stale`; it releases read-only preparations and waits until every workspace read, write, export and Runtime preparation has sent its result. It closes the Session's Executors and, when asked, removes the native home; only then does it reply `released` or `home_removed`. Unfinished cleanup replies `failed` with `cleanup_unconfirmed`, and a retry at the same epoch repeats it. A Runtime that declares `home_removal` unsupported answers `remove_home: true` with `unsupported_operation`, and Core asks it only to release. Core records the release as applied from a matching `released` or `home_removed`, or at once when no Runtime is left to act on it: a release to a Runtime without authority is settled when recorded, and revoking a Runtime settles its releases. Core resends every unacknowledged release to a Runtime when it connects; a release that fails backs off, and the release due longest goes first, so failing releases cannot delay the rest. A quiesced Runtime admits only a release and the matching `environment_resume`, which carries the assignment that quiesced it. +A Session's first bind resolves its Environment owner, which holds the Environment's resources and performs every effect on them; it does not change afterwards. The owner checks each `execution_prepare` configuration against the Environment, including the read-only profile, and fills the installed capabilities before an Executor starts. It applies `runtime_prepare`, lists directories for `workspace_read`, writes files for `workspace_write` and exports outputs for `workspace_export`. The Runtime's dispatcher keeps admission, transfer framing and fencing, and never substitutes another implementation. A self-hosted Runtime's owner is its bound local workspace, which outlives each assignment; the Runtime rejects the bind of any other Session with `assignment_conflict`. A Session without an owner supports none of these operations, and the Runtime rejects each with its typed code: `unsupported_read_preparation` for a read-only preparation, `invalid_configuration` for an Executor configuration with a `local_environment` (except on an agent host, whose Executor binds its Environment itself), `runtime_preparation_unsupported` for `runtime_prepare`, `write_unsupported` for `workspace_write`, and `read_unsupported` for `workspace_read` and `workspace_export`. + +Core records a release and advances the epoch before it sends anything, which withdraws the assignment's attach grant; it then has the relay revoke the Environment's Link resource at its current generation, so the attachments opened under the grant close before the Runtime receives the release. Deleting a Session releases its assignment with `remove_home: true`; releasing its Environment sends `false`. A deletion never revokes a shared Runtime credential. `assignment_release` fences the assignment at once. The Runtime then stops the Session's work: a transfer still receiving its body, or committed but not yet applied, ends with `assignment_stale`; it releases read-only preparations and waits until every workspace read, write, export and Runtime preparation has sent its result. It closes the Session's Executors, then releases what its Environment owner holds and, when asked, removes the native home; only then does it reply `released` or `home_removed`. Unfinished cleanup replies `failed` with `cleanup_unconfirmed`, and a retry at the same epoch repeats it. A Runtime that declares `home_removal` unsupported answers `remove_home: true` with `unsupported_operation`, and Core asks it only to release. Core records the release as applied from a matching `released` or `home_removed`, or at once when no Runtime is left to act on it: a release to a Runtime without authority is settled when recorded, and revoking a Runtime settles its releases. Core resends every unacknowledged release to a Runtime when it connects; a release that fails backs off, and the release due longest goes first, so failing releases cannot delay the rest. A quiesced Runtime admits only a release and the matching `environment_resume`, which carries the assignment that quiesced it. The Runtime answers a Core frame it cannot route with `protocol_error`, which echoes the request's ID and carries its type and an error code. @@ -196,7 +198,7 @@ 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 `workspace_read_preparation` 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 Runtime serves it from its bound local workspace 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 that needs no running Turn uses the read-only preparation profile: `execution_prepare` with `workspace_read_only`, which requires the `workspace_read_preparation` 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. @@ -204,7 +206,7 @@ Request payloads are bounded at 8 KiB and correlation IDs at 128 bytes before ad 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. -`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 execution while it receives or applies a write; a malformed, incomplete or expired transfer never reaches the installer. An exact commit or rejection receipt releases the mutation owner. A missing or ambiguous receipt keeps the uncertainty: 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. Every platform uses the daemon's Go implementation for directory listing, file creation and output export, with no external helper or staging directory. +`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 execution while it receives or applies a write; a malformed, incomplete or expired transfer never reaches the installer. An exact commit or rejection receipt releases the mutation owner. A missing or ambiguous receipt keeps the uncertainty: 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. On every platform the owner lists directories, creates files and exports outputs in the daemon, with no external helper or staging directory. ## MCP connection authority diff --git a/docs/zh/runtime-protocol.md b/docs/zh/runtime-protocol.md index 80cd9a6b1..d27f71134 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: bda8cb83cb7a350d613c44e9b09e6ab78f604643cb31c4f8bfed9e6bcdd45820 +source_hash: 58f0ccc44daf024d3f8079e0f09b78dfefb2e55e558f5663c2add01ee153061a --- 此协议在 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)。 @@ -30,7 +30,7 @@ wire 版本为 [`proto.Version`](https://github.com/MiniMax-AI/OpenAgentCore/blo wire 上每个字段都是 JSON boolean,所有字段都必须出现,包括 `false`。不完整声明编码失败。解码拒绝省略、null、无效和未知字段,以及缺失的 capability 对象。无效 heartbeat 会清空连接的 admission snapshot 并关闭 transport;这不证明原生完成或取消结果。 -每个已准入的 Executor 和 Turn 保留准入时的声明。后续 heartbeat 不能给已有 owner 增加操作。可选操作在任何原生调用前检查此快照;存在 Go interface 不代表支持。已声明操作返回 `agent.ErrUnsupportedOperation` 属于契约违规,与不可用、原生调用失败或不确定写入不同。不确定操作保留回执与所有权,绝不自动重放。工作区支持包括公共 Runtime 工作区实现,因此原生 adapter 不支持的 workspace 方法不会禁用该组合。 +每个已准入的 Executor 和 Turn 保留准入时的声明。后续 heartbeat 不能给已有 owner 增加操作。可选操作在任何原生调用前检查此快照;存在 Go interface 不代表支持。已声明操作返回 `agent.ErrUnsupportedOperation` 属于契约违规,与不可用、原生调用失败或不确定写入不同。不确定操作保留回执与所有权,绝不自动重放。Runtime 根据其 [Environment owner](#session-assignments) 声明 `workspace_read_preparation` 和 `workspace_output_export`,从不由 Harness adapter 声明。 新增字段要求每个生产声明都作出明确决定。契约测试为注册、wire 往返和持久化 boolean 投影逐一枚举字段;共享测试 fixture 单独列出字段,不为未来字段提供默认值。[Harness 接入](../../contracts/agents-api/zh/harness-onboarding.md)负责各声明的 adapter 侧规则。 @@ -118,7 +118,9 @@ Usage frame 和最终 usage snapshot 都携带当前执行的累计测量,替 在一条连接上执行 Session 的第一个操作之前,包括没有 Turn 的 Environment 初始化和文件操作,Core 发送带 Session 的 Environment ID 的 `assignment_bind`,并等待 `assignment_status` `bound`。当 Runtime 是 agent host 且 Environment 有存活的 [Link](./sandbox-link-protocol.md) resource 时,绑定还携带 `resource` 和 `attach_grant`:前者是该 resource,形式与[引导输入](./sandbox-bootstrap.md#launch-input)中的相同;后者是 base64 编码的 grant,agent host 凭它在此分配和 epoch 下打开该 resource generation 上的服务。grant 是机密。Core 不向其他任何 Runtime 发送这两个字段。只带其中一个字段、或带其他 Environment 的 resource 的绑定以 `invalid_request` 失败。以相同的 Environment、resource 和 grant 重复绑定同一分配仍得到 `bound`;该分配的其他绑定以 `assignment_conflict` 失败。Runtime 只在其已绑定的分配下准入 Session frame:较旧的 epoch 或已释放的分配以 `assignment_stale` 失败;其他分配、Session 或 Environment 以 `assignment_conflict` 失败。已启动 Run 的 frame,包括其取消回执,在释放前仍可在启动它的分配下准入。Runtime 已记录回执的重复函数结果或决策只在应用它的分配下得到回答;其他分配以 `assignment_conflict` 失败。 -Core 先记录释放并推进 epoch,再发送任何消息;记录即撤回该分配的 attach grant。随后 Core 让 relay 吊销 Environment 的 Link resource 的当前 generation,使凭该 grant 打开的 attachment 在 Runtime 收到释放之前关闭。删除 Session 以 `remove_home: true` 释放其分配;释放其 Environment 发送 `false`。删除从不吊销共享的 Runtime 凭据。`assignment_release` 立即约束该分配。随后 Runtime 停止 Session 的工作:仍在接收内容、或已提交但尚未应用的传输以 `assignment_stale` 结束;它释放只读准备,并等待每个 workspace 读取、写入、导出和 Runtime 准备发送结果。它关闭 Session 的 Executor,并在要求时删除原生 home;此后才回复 `released` 或 `home_removed`。未完成的清理回复 `failed` 和 `cleanup_unconfirmed`,同一 epoch 的重试会重复清理。声明 `home_removal` 不支持的 Runtime 以 `unsupported_operation` 回答 `remove_home: true`,Core 只要求它释放。Core 根据匹配的 `released` 或 `home_removed` 记录释放已应用;没有 Runtime 能处理该释放时立即记录:发给无授权 Runtime 的释放在记录时即结清,吊销 Runtime 会结清它的释放。Core 在 Runtime 连接时重发所有未确认的释放;失败的释放退避重试,等待最久的释放先发送,因此失败的释放不会拖延其他释放。已 quiesce 的 Runtime 只准入释放和匹配的 `environment_resume`,后者携带使其 quiesce 的分配。 +Session 的第一次绑定确定其 Environment owner,此后不再改变;owner 持有 Environment 的资源,并执行对这些资源的每个作用。owner 根据 Environment 检查每个 `execution_prepare` 配置(包括只读 profile),并在 Executor 启动前填入已安装的能力。它应用 `runtime_prepare`,为 `workspace_read` 列举目录,为 `workspace_write` 写入文件,为 `workspace_export` 导出输出。Runtime 的 dispatcher 保留准入、传输分帧和 fencing,从不替换为其他实现。self-hosted Runtime 的 owner 是其绑定的本地工作区,该工作区比每个分配存续得更久;Runtime 以 `assignment_conflict` 拒绝任何其他 Session 的绑定。没有 owner 的 Session 不支持上述任何操作,Runtime 以各自的类型化错误码拒绝:只读 preparation 为 `unsupported_read_preparation`;带 `local_environment` 的 Executor 配置为 `invalid_configuration`(agent host 除外,其 Executor 自行绑定 Environment);`runtime_prepare` 为 `runtime_preparation_unsupported`;`workspace_write` 为 `write_unsupported`;`workspace_read` 和 `workspace_export` 为 `read_unsupported`。 + +Core 先记录释放并推进 epoch,再发送任何消息;记录即撤回该分配的 attach grant。随后 Core 让 relay 吊销 Environment 的 Link resource 的当前 generation,使凭该 grant 打开的 attachment 在 Runtime 收到释放之前关闭。删除 Session 以 `remove_home: true` 释放其分配;释放其 Environment 发送 `false`。删除从不吊销共享的 Runtime 凭据。`assignment_release` 立即约束该分配。随后 Runtime 停止 Session 的工作:仍在接收内容、或已提交但尚未应用的传输以 `assignment_stale` 结束;它释放只读准备,并等待每个 workspace 读取、写入、导出和 Runtime 准备发送结果。它关闭 Session 的 Executor,随后释放其 Environment owner 持有的资源,并在要求时删除原生 home;此后才回复 `released` 或 `home_removed`。未完成的清理回复 `failed` 和 `cleanup_unconfirmed`,同一 epoch 的重试会重复清理。声明 `home_removal` 不支持的 Runtime 以 `unsupported_operation` 回答 `remove_home: true`,Core 只要求它释放。Core 根据匹配的 `released` 或 `home_removed` 记录释放已应用;没有 Runtime 能处理该释放时立即记录:发给无授权 Runtime 的释放在记录时即结清,吊销 Runtime 会结清它的释放。Core 在 Runtime 连接时重发所有未确认的释放;失败的释放退避重试,等待最久的释放先发送,因此失败的释放不会拖延其他释放。已 quiesce 的 Runtime 只准入释放和匹配的 `environment_resume`,后者携带使其 quiesce 的分配。 Runtime 对无法路由的 Core frame 回复 `protocol_error`,回显请求 ID,并携带其类型和错误码。 @@ -198,7 +200,7 @@ Core 在 Turn outcome 中将接受的值保存为 `engine_error_code` 和 `engin ## 工作区操作 {#workspace-operations} -无需运行 Turn 的工作区读取使用只读 preparation profile:带 `workspace_read_only` 的 `execution_prepare`,要求 `workspace_read_preparation` 能力。仅接受绑定的 Environment 和 resource 身份;不包含 execution option、model 与 MCP 凭据、原生 Session continuation、model 或 tool 输入,owner 拒绝 `execution_start`。Runtime 从绑定的本地工作区提供读取,不启动 Harness 进程。profile 在 Runtime 放弃该 preparation 的所有权后发布 `released`;旧 status snapshot 不发布成功。release 请求、HTTP 断连或远端 socket 关闭本身都不确认释放。 +无需运行 Turn 的工作区读取使用只读 preparation profile:带 `workspace_read_only` 的 `execution_prepare`,要求 `workspace_read_preparation` 能力。仅接受绑定的 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 清理结算后返回。此层没有快照、递归或分页。 @@ -206,7 +208,7 @@ Core 在 Turn outcome 中将接受的值保存为 `engine_error_code` 和 `engin Core 在 Worker 的 Session 调度预约上运行空闲目录读取,活动执行时针对精确 Run。它在有限时 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 在接收或应用写入时排除执行;格式错误、不完整或到期的 transfer 不会到达 installer。精确的 commit 或拒绝回执释放 mutation owner。缺失或有歧义的回执保留不确定性:observer 取消和本地进程退出不能证明没有改变任何内容。公开准入前,Core 在 Session lock 下持久预约写入,跨重启阻止后继 mutation,直到精确结算;请求不重放。各平台都使用 daemon 的 Go 实现进行目录列举、文件创建与输出导出,不使用外部 helper 或 staging directory。 +`workspace_write` 在原生 writer 运行前,通过已确认的 64 KiB frame 传输完整且有界的 body,验证声明的 digest,不运行模型。私有 transfer 限制为 50 MiB,与公开 API 在任何 Runtime 工作前检查的 5 MiB decoded inline 限制独立。Runtime 在接收或应用写入时排除执行;格式错误、不完整或到期的 transfer 不会到达 installer。精确的 commit 或拒绝回执释放 mutation owner。缺失或有歧义的回执保留不确定性:observer 取消和本地进程退出不能证明没有改变任何内容。公开准入前,Core 在 Session lock 下持久预约写入,跨重启阻止后继 mutation,直到精确结算;请求不重放。在各平台上,owner 都在 daemon 内列举目录、创建文件和导出输出,不使用外部 helper 或 staging directory。 ## MCP 连接权限 {#mcp-connection-authority} diff --git a/packages/claude-sdk-adapter/README.md b/packages/claude-sdk-adapter/README.md index b0a0dd717..17a0c06bc 100644 --- a/packages/claude-sdk-adapter/README.md +++ b/packages/claude-sdk-adapter/README.md @@ -38,12 +38,6 @@ Each Turn ends with a result/error and `turn_settled`, independently of process `claudesdk.NewExecutorFactory` binds this bridge to `agent.Executor`. Runtime execution uses the Executor registry for both none and workspace configurations. Its owner context spans all Turns; a Turn's caller cannot replace fixed resources. A failed preparation returns its Executor when cleanup remains unconfirmed. Installed runtime checks are cached by package/file identity, while capability and request validation still run for each Executor configuration. -### Directory listing - -The optional private `agent.WorkspaceDirectoryLister` uses the portable `workspace_directory` bridge feature. Node reads directory metadata under the selected workspace; there is no Linux `/proc/self/fd` dependency. Public daemon Files operations share the Go Binding implementation on all platforms. Their relative path and result bounds define the Files API, not Harness permissions. - -Directory requests are bounded to 8 KiB and 1,000 immediate entries, with explicit truncation, literal names, kinds, and sizes only for regular files. They do not promise ordering, snapshots, recursion, or public pagination. Missing and permission errors are returned only from distinguishable filesystem outcomes; unknown results stop the owner. Each operation closes its directory before a successful receipt. The continuous bridge output consumer retains receipts during preparation and across Turns. Caller cancellation detaches observation without cancelling the Run or discarding an admitted waiter, whose original deadline still applies; owner closure stops admission. The lister alone does not enable public Claude Files; the dedicated Runtime integration supplies public placement and ownership. - ### Command observations Private workspace execution requires the packaged `workspace_command_observations` feature and emits the neutral command snapshots. Match root, current-query native Bash call/result identities after input; ignore historical replay, synthetic and child work. Preserve exact command text and the native per-call textual result, including native rendering or truncation. This is final native output, not incremental stdout/stderr or reconstructed interleaving. Native error results are failed; unambiguous structured interruption is incomplete. Missing results close as incomplete after the observation drain; query cancellation does not overwrite an already observed native failure. Do not infer an exit code from rendered text or supply cwd/duration without qualified native fields. Preparation alone emits no command. Cold continuation must not reissue historical observations. This private translation does not enable public workspace admission, Read/Edit Items or Files ownership. diff --git a/packages/claude-sdk-adapter/src/adapter.ts b/packages/claude-sdk-adapter/src/adapter.ts index a1e8ccd75..d4c655d54 100644 --- a/packages/claude-sdk-adapter/src/adapter.ts +++ b/packages/claude-sdk-adapter/src/adapter.ts @@ -4,7 +4,6 @@ import { StructuredOutput } from "./structured_output.js"; import { toolSearchEnvironment } from "./tool_search.js"; import { Subagents } from "./subagents.js"; import type { Fact } from "./subagent_history.js"; -import { WorkspaceDirectories, type WorkspaceDirectoryEvent } from "./workspace_directories.js"; import { getSessionInfo, query, startup, type McpServerConfig, type Options, type WarmQuery } from "@anthropic-ai/claude-agent-sdk"; import { Inputs, type InputEvent } from "./inputs.js"; import { executorResultUsage, resultUsage, type NativeUsage } from "./usage.js"; @@ -23,7 +22,6 @@ export { parseStart, type Start } from "./request.js"; export type Event = | ExecutorEvent | Fact - | WorkspaceDirectoryEvent | MessageEvent | InputEvent | FunctionEvent @@ -35,7 +33,7 @@ export type Event = | { type: "result"; session_id: string; text: string } | { type: "error"; code: "invalid_request" | "history_unavailable" | "execution_failed" | "cancelled"; engine_error_code?: string; session_id?: string; result_id?: string }; -export async function execute(request: Start | Prepare | ExecutorPrepare, emit: (event: Event) => Promise, abort: AbortController, functions = new FunctionBridge(emit), inputs = new Inputs(immediateInput(request)), directories = new WorkspaceDirectories(emit, abort), turns?: ExecutorTurns): Promise { +export async function execute(request: Start | Prepare | ExecutorPrepare, emit: (event: Event) => Promise, abort: AbortController, functions = new FunctionBridge(emit), inputs = new Inputs(immediateInput(request)), turns?: ExecutorTurns): Promise { const ownerEmit=emit; if(turns) emit=event=>turns.emit(event); const nativeFailure = new NativeFailure(); @@ -135,7 +133,6 @@ export async function execute(request: Start | Prepare | ExecutorPrepare, emit: if ((hooksRequested && initialized.hooks_applied !== true) || children.length !== 1 || !nativeAlive || abort.signal.aborted) { throw new Error("preparation unavailable"); } - if (workspace && process.platform === "linux") await directories.bind(request.cwd); if (declarations?.some(server => "required" in server && server.required)) profile?.verifyRequired(await stream.mcpServerStatus()); if(turns) { turns.configure(stream,(input,output)=>{ @@ -227,7 +224,6 @@ export async function execute(request: Start | Prepare | ExecutorPrepare, emit: inputs.close(); turns?.close(); functions.close(); - try { await directories.close(); } catch { failed = true; classifiedFailure = undefined; } stream?.close(); warm?.close(); abort.signal.removeEventListener("abort", closeInputs); diff --git a/packages/claude-sdk-adapter/src/main.ts b/packages/claude-sdk-adapter/src/main.ts index e3ba0bd11..891e77080 100644 --- a/packages/claude-sdk-adapter/src/main.ts +++ b/packages/claude-sdk-adapter/src/main.ts @@ -1,5 +1,4 @@ import { ExecutorTurns } from "./executor_protocol.js"; -import { WorkspaceDirectories } from "./workspace_directories.js"; import { Inputs } from "./inputs.js"; import { FunctionBridge } from "./function_bridge.js"; import { createInterface } from "node:readline"; @@ -27,7 +26,6 @@ try { await emit(invalid && (event.type === "error" || event.type === "result") ? { type: "error", code: "invalid_request" } : event); }; const functions = new FunctionBridge(output); - const directories = new WorkspaceDirectories(output, abort); const prompts = new Inputs(immediateInput(request)); const turns=request.type==="executor_prepare" ? new ExecutorTurns(emit,abort) : undefined; const incoming = (async () => { @@ -35,9 +33,7 @@ try { for await (const line of { [Symbol.asyncIterator]: () => input }) { if (Buffer.byteLength(line) > 1024 * 1024) throw new Error("Invalid input."); const value: unknown = JSON.parse(line); - if (value && typeof value === "object" && "type" in value && value.type === "workspace_directory") { - directories.submit(value as Record); - } else if(turns) { + if(turns) { if(!value || typeof value!=="object" || Array.isArray(value)) throw new Error("invalid_request"); const control=value as Record; if(control.type==="turn_start") await turns.start(control); @@ -54,7 +50,7 @@ try { } catch { invalid = request.type === "prepare"; abort.abort(); } })(); - try { await execute(request, output, abort, functions, prompts, directories, turns); } + try { await execute(request, output, abort, functions, prompts, turns); } catch { await output({ type: "error", code: "execution_failed" }); } finally { prompts.close(); diff --git a/packages/claude-sdk-adapter/src/runtime_check.ts b/packages/claude-sdk-adapter/src/runtime_check.ts index b8a4a1852..fe093b1fd 100644 --- a/packages/claude-sdk-adapter/src/runtime_check.ts +++ b/packages/claude-sdk-adapter/src/runtime_check.ts @@ -45,7 +45,7 @@ try { assert.equal(smoke.error, undefined, "bridge_unavailable"); assert.equal(smoke.status, 0, "bridge_unavailable"); assert.deepEqual(JSON.parse(smoke.stdout), { type: "error", code: "invalid_request" }); - process.stdout.write(JSON.stringify({ type: "runtime_ready", protocol: 3, features: ["executor_reuse", ...(["linux", "darwin", "win32"].includes(process.platform) ? ["workspace_directory", "local_runtime_v2", "workspace_functions", "workspace_structured_output", "workspace_tool_search", "workspace_mcp_http"] : []), "message_images", "function_result_images", "tool_search", "structured_output", "subagent_resources", "mcp_http_tools", "mcp_http_bearer_auth", "mcp_http_required", "workspace_tools", "workspace_prepare", "workspace_command_observations"], node: process.versions.node, sdk: sdk.version, mcp: mcp.version, native: nativeVersion, native_path: relative(root, binary) }) + "\n"); + process.stdout.write(JSON.stringify({ type: "runtime_ready", protocol: 3, features: ["executor_reuse", ...(["linux", "darwin", "win32"].includes(process.platform) ? ["local_runtime_v2", "workspace_functions", "workspace_structured_output", "workspace_tool_search", "workspace_mcp_http"] : []), "message_images", "function_result_images", "tool_search", "structured_output", "subagent_resources", "mcp_http_tools", "mcp_http_bearer_auth", "mcp_http_required", "workspace_tools", "workspace_prepare", "workspace_command_observations"], node: process.versions.node, sdk: sdk.version, mcp: mcp.version, native: nativeVersion, native_path: relative(root, binary) }) + "\n"); } catch { // Native diagnostics can include operator environment; never forward them. process.stdout.write(JSON.stringify({ type: "runtime_unavailable" }) + "\n"); diff --git a/packages/claude-sdk-adapter/src/workspace_directories.ts b/packages/claude-sdk-adapter/src/workspace_directories.ts deleted file mode 100644 index 355bf0c0a..000000000 --- a/packages/claude-sdk-adapter/src/workspace_directories.ts +++ /dev/null @@ -1,75 +0,0 @@ -import { join, isAbsolute } from "node:path"; -import { lstat, opendir, stat } from "node:fs/promises"; - -export type WorkspaceDirectoryEntry = { name: string; kind: "file" | "directory" | "symlink" | "other"; size_bytes?: number }; -export type WorkspaceDirectoryEvent = { type: "workspace_directory"; id: string } & ( - { entries: WorkspaceDirectoryEntry[]; truncated: boolean } | - { error: "invalid" | "busy" | "unavailable" | "uncertain" | "not_found" | "permission" } -); -export class WorkspaceDirectories { - private root?: string; - private stopped = false; - private pending?: Promise; - - constructor(private readonly emit: (event: WorkspaceDirectoryEvent) => Promise, private readonly abort: AbortController) {} - - async bind(root: string): Promise { - if (this.root || this.stopped || !isAbsolute(root) || !(await stat(root)).isDirectory()) throw new Error("directory listing unavailable"); - this.root = root; - } - - submit(value: Record): void { - if (typeof value.id !== "string" || !/^[a-zA-Z0-9-]{1,128}$/.test(value.id)) throw new Error("invalid_request"); - const id = value.id; - let error: "invalid" | "busy" | "unavailable" | undefined; - if (Buffer.byteLength(JSON.stringify(value)) > 8192 || Object.keys(value).some(key => !["type", "id", "directory", "max_entries"].includes(key)) || - typeof value.directory !== "string" || /[\x00\\\r\n]/.test(value.directory) || - (value.directory !== "" && value.directory.split("/").some(part => !part || part === "." || part === "..")) || - !Number.isInteger(value.max_entries) || (value.max_entries as number) < 1 || (value.max_entries as number) > 1000) error = "invalid"; - else if (this.stopped || this.abort.signal.aborted || !this.root) error = "unavailable"; - else if (this.pending) error = "busy"; - if (error) { void this.emit({ type: "workspace_directory", id, error }).catch(() => this.abort.abort()); return; } - this.pending = this.list(id, value.directory as string, value.max_entries as number).finally(() => { this.pending = undefined; }); - } - - private async list(id: string, directory: string, maxEntries: number): Promise { - try { - const result = await this.enumerate(directory, maxEntries); - await this.emit({ type: "workspace_directory", id, ...result }); - } catch (error) { - const code = (error as NodeJS.ErrnoException).code; - const classified = code === "ENOENT" ? "not_found" : code === "ENOTDIR" ? "invalid" : - ["EACCES", "EPERM", "ELOOP"].includes(code ?? "") ? "permission" : "uncertain"; - await this.emit({ type: "workspace_directory", id, error: classified }).catch(() => this.abort.abort()); - if (classified === "uncertain") { this.stopped = true; this.abort.abort(); } - } - } - - private async enumerate(directory: string, maxEntries: number): Promise<{ entries: WorkspaceDirectoryEntry[]; truncated: boolean }> { - const path = join(this.root!, ...directory.split("/")); - const dir = await opendir(path, { encoding: "buffer" as BufferEncoding }); - try { - const entries: WorkspaceDirectoryEntry[] = []; - while (true) { - if (this.abort.signal.aborted) throw new Error("owner closed"); - const entry = await dir.read(); - if (!entry) return { entries, truncated: false }; - if (entries.length === maxEntries) return { entries, truncated: true }; - const rawName: unknown = entry.name; - if (!Buffer.isBuffer(rawName)) throw new Error("invalid entry name"); - const name = rawName.toString("utf8"); - if (!Buffer.from(name).equals(rawName) || !name || /[\x00/]/.test(name) || name === "." || name === "..") throw new Error("unsupported entry name"); - const stat = await lstat(join(path,name)); - const kind = stat.isFile() ? "file" : stat.isDirectory() ? "directory" : stat.isSymbolicLink() ? "symlink" : "other"; - if (kind === "file" && (!Number.isSafeInteger(stat.size) || stat.size < 0)) throw new Error("invalid file size"); - entries.push({ name, kind, ...(kind === "file" ? { size_bytes: stat.size } : {}) }); - } - } finally { await dir.close(); } - } - - async close(): Promise { - this.stopped = true; - await this.pending; - this.root = undefined; - } -} diff --git a/packages/claude-sdk-adapter/tests/workspace_directories.test.mjs b/packages/claude-sdk-adapter/tests/workspace_directories.test.mjs deleted file mode 100644 index d9c3ae1bb..000000000 --- a/packages/claude-sdk-adapter/tests/workspace_directories.test.mjs +++ /dev/null @@ -1,56 +0,0 @@ -import assert from "node:assert/strict"; -import { mkdtemp, mkdir, writeFile, symlink, rm, rename, readdir } from "node:fs/promises"; -import { tmpdir } from "node:os"; -import { join } from "node:path"; -import test from "node:test"; -import { WorkspaceDirectories } from "../dist/workspace_directories.js"; - -async function fixture(t) { - const root = await mkdtemp(join(tmpdir(), "oac-directories-")); - const workspace = join(root, "workspace"); - await mkdir(workspace); - const abort = new AbortController(); - let receipt; - const directories = new WorkspaceDirectories(event => { receipt?.(event); return Promise.resolve(); }, abort); - await directories.bind(workspace); - t.after(async () => { await directories.close(); await rm(root, { recursive: true, force: true }); }); - let sequence = 0; - return { root, workspace, directories, abort, list: async (directory = "", max_entries = 1000) => { - const result = new Promise(resolve => { receipt = resolve; }); - directories.submit({ type: "workspace_directory", id: `list-${++sequence}`, directory, max_entries }); - const value = await result; - await new Promise(resolve => setImmediate(resolve)); - return value; - } }; -} - -test("bounded literal metadata, missing directory, and traversal denial", async t => { - const f = await fixture(t); - await mkdir(join(f.workspace, "nested")); - await writeFile(join(f.workspace, "binary"), Buffer.from([0, 255, 0])); - await writeFile(join(f.workspace, "empty"), ""); - await writeFile(join(f.root, "protected"), "secret"); - await symlink(f.root, join(f.workspace, "outside")); - await symlink("nested", join(f.workspace, "inside")); - const value = await f.list(); - assert.equal(value.truncated, false); - assert.deepEqual(value.entries.sort((a, b) => a.name.localeCompare(b.name)), [ - { name: "binary", kind: "file", size_bytes: 3 }, { name: "empty", kind: "file", size_bytes: 0 }, - { name: "inside", kind: "symlink" }, { name: "nested", kind: "directory" }, { name: "outside", kind: "symlink" }, - ]); - assert.deepEqual((await f.list("nested")).entries, []); - assert.equal((await f.list("", 1)).truncated, true); - assert.equal((await f.list("absent")).error, "not_found"); - assert.deepEqual((await f.list("inside")).entries, []); - for (const path of ["/", "a//b", ".", "..", "a/../b", "a\\b", "a\n"]) assert.equal((await f.list(path)).error, "invalid"); - for (const limit of [0, -1, 1001, 1.1]) assert.equal((await f.list("", limit)).error, "invalid"); -}); - -test("literal Unicode names are preserved and undecodable names fail explicitly", {skip:process.platform === "win32"}, async t => { - const f = await fixture(t); - await writeFile(join(f.workspace, "字\ufffd"), "ok"); - assert.deepEqual((await f.list()).entries, [{ name: "字\ufffd", kind: "file", size_bytes: 2 }]); - await writeFile(Buffer.concat([Buffer.from(f.workspace + "/"), Buffer.from([255])]), "invalid"); - assert.equal((await f.list()).error, "uncertain"); - assert.equal(f.abort.signal.aborted, true); -}); diff --git a/scripts/check-claude-sdk-runtime.mjs b/scripts/check-claude-sdk-runtime.mjs index c6aee4f11..16e079858 100644 --- a/scripts/check-claude-sdk-runtime.mjs +++ b/scripts/check-claude-sdk-runtime.mjs @@ -23,7 +23,7 @@ assert.equal(probe.status, 0, "Exported runtime is unavailable"); const report = JSON.parse(probe.stdout); assert.equal(report.type, "runtime_ready"); assert.equal(report.protocol, 3); -assert.deepEqual(report.features, ["executor_reuse", ...(["linux", "darwin", "win32"].includes(process.platform) ? ["workspace_directory", "local_runtime_v2", "workspace_functions", "workspace_structured_output", "workspace_tool_search", "workspace_mcp_http"] : []), "message_images", "function_result_images", "tool_search", "structured_output", "subagent_resources", "mcp_http_tools", "mcp_http_bearer_auth", "mcp_http_required", "workspace_tools", "workspace_prepare", "workspace_command_observations"]); +assert.deepEqual(report.features, ["executor_reuse", ...(["linux", "darwin", "win32"].includes(process.platform) ? ["local_runtime_v2", "workspace_functions", "workspace_structured_output", "workspace_tool_search", "workspace_mcp_http"] : []), "message_images", "function_result_images", "tool_search", "structured_output", "subagent_resources", "mcp_http_tools", "mcp_http_bearer_auth", "mcp_http_required", "workspace_tools", "workspace_prepare", "workspace_command_observations"]); assert.equal(report.sdk, source.dependencies["@anthropic-ai/claude-agent-sdk"]); assert.equal(report.mcp, source.dependencies["@modelcontextprotocol/sdk"]); console.log(`Verified exported SDK ${report.sdk}, MCP ${report.mcp}, ${report.native}`); diff --git a/scripts/name-allowlist.json b/scripts/name-allowlist.json index 768db4786..474fd61a1 100644 --- a/scripts/name-allowlist.json +++ b/scripts/name-allowlist.json @@ -89,11 +89,6 @@ "regex": "\"agents-api-(?:session)?\"", "reason": "AgentStateKey is a persisted daemon/native-session resume identity carried on the existing machine wire contract; changing it would select another session." }, - { - "path": "apps/daemon/internal/localworkspace/write.go", - "regex": "\"agents-api-(?:session)?\"", - "reason": "AgentStateKey is a persisted daemon/native-session resume identity carried on the existing machine wire contract; changing it would select another session." - }, { "path": "services/core/internal/execution/directory_preparation.go", "regex": "\"agents-api-(?:session)?\"",