Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 0 additions & 4 deletions apps/daemon/internal/dispatch/prepared_handoff.go
Original file line number Diff line number Diff line change
Expand Up @@ -89,10 +89,6 @@ func (r *Router) preparedOperationLocked(state *sessionState) (agent.Turn, func(
return state.session, handoff.operations.RUnlock, true
}

func (r *Router) runRouteOpenLocked(state *sessionState) bool {
return state != nil && !r.closed && state.ctx.Err() == nil && !state.steeringClosed && state.preparedHandoff.release == nil
}

// claimPreparedReleaseLocked closes admission permanently and returns the
// current native attempt. retry starts a new serialized attempt only after the
// previous one failed. Router.mu must be held.
Expand Down
18 changes: 16 additions & 2 deletions apps/daemon/internal/dispatch/prepared_handoff_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -180,15 +180,29 @@ func TestPreparedHandoffDuplicateStartDoesNotReexecuteDuringPublication(t *testi
if ack := lastSteeringAck(t, sender.recSender, "run", "publication-steering"); ack.ErrorCode != "not_ready" {
t.Fatalf("pre-publication steering = %+v", ack)
}
read := proto.WorkspaceReadPayload{RunID: "run", EnvironmentID: preparationEnvironmentID, MaxEntries: 1}
read := map[string]any{"run_id": "run", "environment_id": preparationEnvironmentID, "max_entries": 1}
if err := r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceRead, "publication-read", read)); err != nil {
t.Fatal(err)
}
var readResult proto.WorkspaceReadResultPayload
frames := sender.snapshot()
if len(frames) == 0 || frames[len(frames)-1].Type != proto.TypeWorkspaceReadResult || frames[len(frames)-1].DecodePayload(&readResult) != nil || readResult.ErrorCode != "resource_unavailable" {
if len(frames) == 0 || frames[len(frames)-1].Type != proto.TypeWorkspaceReadResult || frames[len(frames)-1].DecodePayload(&readResult) != nil || readResult.ErrorCode != "invalid_request" {
t.Fatalf("pre-publication workspace read = %+v", frames)
}
readPreparation := preparationRequest()
readPreparation.Configuration = proto.PromptRequestPayload{AgentKind: "prepared", LocalEnvironment: &proto.LocalEnvironment{ID: preparationEnvironmentID}, WorkspaceReadOnly: true}
if err := r.Handle(t.Context(), mustEnv(t, proto.TypeExecutionPrepare, "publication-read-prepare", readPreparation)); err != nil {
t.Fatal(err)
}
readReady := waitPreparationStatus(t, sender.recSender, "publication-read-prepare", "ready", "")
if err := r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceRead, "publication-independent-read", proto.WorkspaceReadPayload{Handle: readReady.Handle, EnvironmentID: preparationEnvironmentID, MaxEntries: 1})); err != nil {
t.Fatal(err)
}
if got := waitWorkspaceRead(t, sender.recSender, "publication-independent-read"); got.Outcome != "completed" || !got.CloseAcknowledged || got.Directory == nil {
t.Fatal("read-only preparation depends on started publication", got)
}
_ = r.Handle(t.Context(), mustEnv(t, proto.TypeExecutionRelease, "publication-read-prepare", proto.ExecutionReleasePayload{Handle: readReady.Handle}))
waitPreparationStatus(t, sender.recSender, "publication-read-prepare", "released", "")
if session.functions.Load() != 0 || session.steers.Load() != 0 {
t.Fatal("private Session accepted work before started publication")
}
Expand Down
35 changes: 25 additions & 10 deletions apps/daemon/internal/dispatch/workspace_directory_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ import (
"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto"
)

func TestWorkspaceDirectoryRetainsEnvironmentAndTransferredOwner(t *testing.T) {
func TestWorkspaceDirectoryUsesReadOnlyOwnerBesideNativeExecution(t *testing.T) {
sender := &recSender{}
p := &controlledPreparation{closed: make(chan struct{})}
p.start = func(ctx context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) {
Expand All @@ -21,7 +21,13 @@ func TestWorkspaceDirectoryRetainsEnvironmentAndTransferredOwner(t *testing.T) {
r := ownedPreparationRouter(t, sender, time.Minute, func(context.Context, proto.PromptRequestPayload) (preparedFixture, error) { return p, nil }, owner)
_ = r.Handle(t.Context(), mustEnv(t, proto.TypeExecutionPrepare, "prepare", preparationRequest()))
ready := waitPreparationStatus(t, sender, "prepare", "ready", "")
request := proto.WorkspaceReadPayload{Handle: ready.Handle, EnvironmentID: preparationEnvironmentID, MaxEntries: 1}
configuration := preparationRequest()
configuration.Configuration = proto.PromptRequestPayload{AgentKind: "prepared", LocalEnvironment: &proto.LocalEnvironment{ID: preparationEnvironmentID}, WorkspaceReadOnly: true}
if err := r.Handle(t.Context(), mustEnv(t, proto.TypeExecutionPrepare, "read-prepare", configuration)); err != nil {
t.Fatal(err)
}
readReady := waitPreparationStatus(t, sender, "read-prepare", "ready", "")
request := proto.WorkspaceReadPayload{Handle: readReady.Handle, EnvironmentID: preparationEnvironmentID, MaxEntries: 1}
for _, phase := range []string{"idle", "active"} {
_ = r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceRead, phase, request))
result := waitWorkspaceRead(t, sender, phase)
Expand All @@ -34,24 +40,33 @@ func TestWorkspaceDirectoryRetainsEnvironmentAndTransferredOwner(t *testing.T) {
if got := waitWorkspaceRead(t, sender, phase+"-foreign"); got.ErrorCode != proto.AssignmentConflict {
t.Fatal(got)
}
native := request
native.Handle = ready.Handle
_ = r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceRead, phase+"-native", native))
if got := waitWorkspaceRead(t, sender, phase+"-native"); got.ErrorCode != "resource_unavailable" {
t.Fatal("execution preparation admitted a File read", got)
}
if phase == "idle" {
_ = r.Handle(t.Context(), mustEnv(t, proto.TypeExecutionStart, "prepare", proto.ExecutionStartPayload{Handle: ready.Handle, ExecutorID: ready.ExecutorID, RunID: "run", Input: proto.TextInput("start")}))
waitPreparationStatus(t, sender, "prepare", "started", "")
_ = r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceRead, "stale", request))
if got := waitWorkspaceRead(t, sender, "stale"); got.Outcome != "rejected" {
t.Fatal(got)
}
request.Handle, request.RunID = "", "run"
}
}
for index, bad := range []proto.WorkspaceReadPayload{
{RunID: "run", EnvironmentID: preparationEnvironmentID, MaxEntries: proto.WorkspaceDirectoryMaxEntries + 1},
{Handle: ready.Handle, RunID: "run", EnvironmentID: preparationEnvironmentID, MaxEntries: 2},
for index, bad := range []any{
proto.WorkspaceReadPayload{Handle: readReady.Handle, EnvironmentID: preparationEnvironmentID, MaxEntries: proto.WorkspaceDirectoryMaxEntries + 1},
proto.WorkspaceReadPayload{EnvironmentID: preparationEnvironmentID, MaxEntries: 2},
map[string]any{"run_id": "run", "environment_id": preparationEnvironmentID, "max_entries": 2},
map[string]any{"handle": readReady.Handle, "run_id": "run", "environment_id": preparationEnvironmentID, "max_entries": 2},
} {
id := fmt.Sprintf("invalid-%d", index)
_ = r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceRead, id, bad))
if got := waitWorkspaceRead(t, sender, id); got.Outcome != "rejected" || got.ErrorCode != "invalid_request" || got.Directory != nil {
t.Fatal("malformed directory control reached a resource", got)
}
}
_ = r.Handle(t.Context(), mustEnv(t, proto.TypeExecutionRelease, "read-prepare", proto.ExecutionReleasePayload{Handle: readReady.Handle}))
waitPreparationStatus(t, sender, "read-prepare", "released", "")
_ = r.Handle(t.Context(), mustEnv(t, proto.TypeWorkspaceRead, "released-read", request))
if got := waitWorkspaceRead(t, sender, "released-read"); got.ErrorCode != "resource_unavailable" {
t.Fatal("released read owner remained usable", got)
}
}
2 changes: 1 addition & 1 deletion apps/daemon/internal/dispatch/workspace_export.go
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@ func (r *Router) handleWorkspaceExport(ctx context.Context, env proto.Envelope)
}
environment, code := r.workspaceResourceLocked(env.Assignment, proto.WorkspaceReadPayload{Handle: request.Handle, EnvironmentID: request.EnvironmentID})
p := r.preparations[request.Handle]
if code == "" && (u != nil || r.workspaceWrites[env.Assignment.SessionID] != nil || p.executor != nil) {
if code == "" && (u != nil || r.workspaceWrites[env.Assignment.SessionID] != nil) {
code = "resource_unavailable"
}
if code != "" {
Expand Down
16 changes: 4 additions & 12 deletions apps/daemon/internal/dispatch/workspace_read.go
Original file line number Diff line number Diff line change
Expand Up @@ -61,7 +61,7 @@ func (r *Router) handleWorkspaceRead(ctx context.Context, env proto.Envelope) er
}

// workspaceResourceLocked returns the Environment owner of ref's Session once
// ref admits the read's ready preparation handle or open Run.
// ref admits the read-only preparation's ready handle.
func (r *Router) workspaceResourceLocked(ref proto.AssignmentRef, request proto.WorkspaceReadPayload) (Environment, string) {
if code := r.admitLocked(ref, ref.SessionID, request.EnvironmentID); code != "" {
return nil, code
Expand All @@ -70,17 +70,9 @@ func (r *Router) workspaceResourceLocked(ref proto.AssignmentRef, request proto.
if environment == nil {
return nil, "read_unsupported"
}
if request.Handle != "" {
p := r.preparations[request.Handle]
if p == nil || p.request.Assignment != ref || p.environmentID != request.EnvironmentID || p.status.State != "ready" ||
!p.owns || p.busy || p.ctx.Err() != nil || !time.Now().Before(p.deadline) {
return nil, "resource_unavailable"
}
return environment, ""
}
s := r.sessions[request.RunID]
if s == nil || s.assignment != ref || s.environmentID != request.EnvironmentID || s.session == nil ||
!r.runRouteOpenLocked(s) {
p := r.preparations[request.Handle]
if p == nil || p.request.Assignment != ref || p.environmentID != request.EnvironmentID || p.status.State != "ready" ||
p.executor != nil || !p.owns || p.busy || p.ctx.Err() != nil || !time.Now().Before(p.deadline) {
return nil, "resource_unavailable"
}
return environment, ""
Expand Down
6 changes: 3 additions & 3 deletions docs/runtime-protocol.md
Original file line number Diff line number Diff line change
Expand Up @@ -194,13 +194,13 @@ Core stores accepted values in the Turn outcome as `engine_error_code` and `engi

## Workspace operations

A workspace read that needs no running Turn uses the read-only preparation profile: `execution_prepare` with `workspace_read_only`, which requires the `local_environment` capability. It accepts only the bound Environment and resource identity; execution options, model and MCP credentials, native Session continuation and model or tool input are excluded, and the owner rejects `execution_start`. The Session's Environment owner serves it without starting a Harness process. The profile publishes `released` after the Runtime drops the preparation's ownership; a stale status snapshot never publishes success. A release request, HTTP disconnect or remote socket closure alone does not confirm the release.
A workspace read uses the read-only preparation profile independently of native preparation or a running Turn: `execution_prepare` with `workspace_read_only`, which requires the `local_environment` capability. It accepts only the bound Environment and resource identity; execution options, model and MCP credentials, native Session continuation and model or tool input are excluded, and the owner rejects `execution_start`. The Session's Environment owner serves it without starting a Harness process. The profile publishes `released` after the Runtime drops the preparation's ownership; a stale status snapshot never publishes success. A release request, HTTP disconnect or remote socket closure alone does not confirm the release.

`workspace_read` lists one workspace-relative directory (an empty path selects the root) of an existing preparation handle, or the Run it was transferred to, on the same authenticated device connection, with the exact frozen Environment identity; callers cannot supply sockets, credentials or workspace roots. A result carries at most `max_entries` (1 to 1024) single-component UTF-8 names of at most 255 bytes each, the entry kind, regular-file sizes and explicit truncation, and is returned only after directory access and handle cleanup settle. There is no snapshot, recursion or pagination at this layer.
`workspace_read` lists one workspace-relative directory (an empty path selects the root) through a ready read-only preparation handle on the same authenticated device connection, with the exact frozen Environment identity; execution preparation handles are rejected, and callers cannot supply Run identities, sockets, credentials or workspace roots. A result carries at most `max_entries` (1 to 1024) single-component UTF-8 names of at most 255 bytes each, the entry kind, regular-file sizes and explicit truncation, and is returned only after directory access and handle cleanup settle. There is no snapshot, recursion or pagination at this layer.

Request payloads are bounded at 8 KiB and correlation IDs at 128 bytes before admission; oversized IDs are not echoed, and oversized trace metadata is omitted from replies. These are private transport limits, not public Files parameters. A safe native rejection carries no directory; an interrupted or ambiguous read stays unknown and stops further reads on that owner. A dispatched read keeps its bounded waiter across observer cancellation and resource transfer or release, and resource closure stops new admission. The gateway bounds subscriptions and never retries or replays a read on reconnect; a duplicate pending operation ID cannot start another read.

Core runs an idle directory read on the Worker's Session scheduling reservation and targets the exact Run during active execution. It keeps the reservation through the bounded read and release, returns data only after a confirmed close (an incomplete read or uncertain cleanup returns unavailable without data), releases the reservation before delivering the result, and revokes the scoped read credential on completion or failure. The Runtime keeps uncertain cleanup ownership and capacity. The [Environment Files contract](../contracts/agents-api/environment-files.md) owns public authorization, paths and pagination.
Core uses an independent read-only preparation for every directory read and reserves the Worker's Session scheduling slot when idle. It keeps that reservation through the bounded read and release, returns data only after a confirmed close (an incomplete read or uncertain cleanup returns unavailable without data), releases the reservation before delivering the result, and revokes the scoped read credential on completion or failure. The Runtime keeps uncertain cleanup ownership and capacity. The [Environment Files contract](../contracts/agents-api/environment-files.md) owns public authorization, paths and pagination.

`workspace_write` transfers a complete bounded body in acknowledged 64 KiB frames before the native writer runs, verifies the declared digest and runs no model. The private transfer bound is 50 MiB, separate from the public 5 MiB decoded inline bound that the API checks before any Runtime work. The Runtime excludes the Session's execution while it receives or applies a write; a malformed, incomplete or expired transfer never reaches the installer. An exact commit or rejection receipt settles the write. A missing or ambiguous receipt leaves the uncertainty with the [Environment owner](#session-assignments): observer cancellation and local process exit cannot prove that nothing changed. Before public admission Core durably reserves the write under the Session lock and blocks successor mutations across restarts until exact settlement; the request is never replayed. The agent host's owner lists directories, creates files and exports outputs over [File access](./file-access-protocol.md) on its Link attachment.

Expand Down
Loading
Loading