From a2e7cf93fe9eb23ca72a057bd99e6e960b417bc4 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Wed, 7 Oct 2026 23:26:06 +0000 Subject: [PATCH 1/3] Scope dispatch transfers, quiesce and supersede to their Session and Environment A bind at a higher epoch of the bound assignment fences the earlier epoch's work as a release does, closes its Executor and owner, and binds. Transfer slots are held per Session under a per-connection memory bound. environment_quiesce and environment_resume apply to the named Environment: its Executors close and its owners release their attachments while other Environments keep running. --- .../agenthost/environment_linux_test.go | 46 ++++ .../internal/agenthost/view_linux_test.go | 7 +- apps/daemon/internal/dispatch/assignment.go | 144 ++++++++--- .../internal/dispatch/assignment_test.go | 57 +++++ apps/daemon/internal/dispatch/cancellation.go | 6 +- apps/daemon/internal/dispatch/executor.go | 4 +- .../dispatch/native_file_results_test.go | 3 +- apps/daemon/internal/dispatch/preparation.go | 6 +- .../internal/dispatch/preparation_start.go | 4 +- .../internal/dispatch/prepared_handoff.go | 3 + apps/daemon/internal/dispatch/router.go | 42 +++- .../internal/dispatch/runtime_preparation.go | 32 +-- .../dispatch/runtime_preparation_test.go | 127 ++++++++-- apps/daemon/internal/dispatch/shutdown.go | 16 +- apps/daemon/internal/dispatch/suspend.go | 231 ++++++++++++++---- apps/daemon/internal/dispatch/suspend_test.go | 111 ++++++--- .../internal/dispatch/workspace_export.go | 14 +- .../dispatch/workspace_export_test.go | 6 +- .../internal/dispatch/workspace_write.go | 20 +- docs/runtime-protocol.md | 10 +- docs/zh/runtime-protocol.md | 12 +- 21 files changed, 704 insertions(+), 197 deletions(-) diff --git a/apps/daemon/internal/agenthost/environment_linux_test.go b/apps/daemon/internal/agenthost/environment_linux_test.go index de271aa69..83fbbee45 100644 --- a/apps/daemon/internal/agenthost/environment_linux_test.go +++ b/apps/daemon/internal/agenthost/environment_linux_test.go @@ -204,6 +204,52 @@ func TestEnvironmentOwnerServesTheSandbox(t *testing.T) { } } +// TestSupersedingBindRebindsTheOwner checks that a bind of the Session's +// assignment at a higher epoch with a new grant drains the owner's +// attachment and gives the owner the new binding, under which its next write +// attaches to the sandbox again. +func TestSupersedingBindRebindsTheOwner(t *testing.T) { + if os.Getenv(gateEnv) != "1" { + t.Skipf("set %s=1 and run the test binary as root in a throwaway container; see the view suite", gateEnv) + } + sb := startSandbox(t, os.Getenv(sandboxIOEnv)) + for _, name := range []string{"before.txt", "after.txt"} { + if err := os.RemoveAll(path.Join(sandboxWorkspace, name)); err != nil { + t.Fatal(err) + } + } + if err := os.MkdirAll(sandboxWorkspace, 0o777); err != nil { + t.Fatal(err) + } + cfg := Config{StateDir: t.TempDir(), RelayURL: sb.url, RuntimeID: sandboxwire.NewID(), Credential: []byte("runtime-credential"), Harnesses: agent.NewRegistry()} + sb.auth.AddRuntime(cfg.Credential, cfg.RuntimeID) + sb.ready(t, cfg) + dm := &daemon{host: &Host{cfg: cfg, owners: owners{d: deps{dial: relayDial(cfg)}}}} + dm.route(t, agent.NewRegistry()) + b := sb.bind(cfg.RuntimeID, time.Minute) + dm.assign(t, b) + if r := dm.write(t, b, "before.txt", []byte("before")); r.Outcome != "completed" || dm.host.drained(t, b) { + t.Fatalf("the write at epoch 1 is %+v", r) + } + + next := b + next.AssignmentEpoch, next.AttachGrant = 2, []byte("grant-"+sandboxwire.NewID().String()) + sb.grant(next, cfg.RuntimeID, time.Minute) + dm.assign(t, next) + o := dm.host.owner(t, next) + rebound, drained := sameBinding(o.binding, next), o.link == nil + o.release() + if !rebound || !drained { + t.Fatalf("after the superseding bind the owner has the new binding %t and no attachment %t", rebound, drained) + } + if r := dm.write(t, next, "after.txt", []byte("after")); r.Outcome != "completed" { + t.Fatalf("the write at epoch 2 is %+v", r) + } + if body, err := os.ReadFile(path.Join(sandboxWorkspace, "after.txt")); err != nil || string(body) != "after" { + t.Fatalf("the sandbox has %q, %v", body, err) + } +} + // TestUnreachableSandboxRejectsRuntimePreparation checks that a // runtime_prepare whose owner cannot reach the sandbox ends rejected, which // leaves the Router free to run another Session's and to shut down. diff --git a/apps/daemon/internal/agenthost/view_linux_test.go b/apps/daemon/internal/agenthost/view_linux_test.go index a8cb478d7..d3dfc3945 100644 --- a/apps/daemon/internal/agenthost/view_linux_test.go +++ b/apps/daemon/internal/agenthost/view_linux_test.go @@ -566,10 +566,15 @@ func (sb *sandbox) stop() { // runtimeID for lease. func (sb *sandbox) bind(runtimeID sandboxwire.ID, lease time.Duration) Binding { b := newBinding(sb.resource) + sb.grant(b, runtimeID, lease) + return b +} + +// grant lets runtimeID attach b's Session to the resource for lease. +func (sb *sandbox) grant(b Binding, runtimeID sandboxwire.ID, lease time.Duration) { sb.auth.AddGrant(b.AttachGrant, sandboxlinktest.Grant{RuntimeID: runtimeID, Resource: b.Resource, SessionID: b.SessionID, AssignmentID: b.AssignmentID, AssignmentEpoch: b.AssignmentEpoch, Lease: lease, Services: []sandboxlink.Service{sandboxlink.ServiceFile, sandboxlink.ServiceProcess, sandboxlink.ServiceNetwork}}) - return b } func (sb *sandbox) dial(t *testing.T, cfg Config) *sandboxlink.AttachLink { diff --git a/apps/daemon/internal/dispatch/assignment.go b/apps/daemon/internal/dispatch/assignment.go index edebffd51..b7c399542 100644 --- a/apps/daemon/internal/dispatch/assignment.go +++ b/apps/daemon/internal/dispatch/assignment.go @@ -22,12 +22,16 @@ type assignmentState struct { // 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. - work sync.WaitGroup - // cleanup serializes release cleanups, so a retried release never closes - // the owner while an earlier Close runs. + // superseding is the bind that superseded the assignment at ref, held + // released until the earlier epoch's work and owner have settled. + superseding *proto.AssignmentBindPayload + // work counts the Session's admitted reads, writes, exports, Runtime + // preparations, Run terminals and cancellation receipts until each has + // sent its terminal result. A release, a superseding bind and a quiesce + // wait for it, and a released assignment admits no more. + work dispatchWork + // cleanup serializes release and supersede cleanups, so a retry never + // closes the owner while an earlier Close runs. cleanup sync.Mutex } @@ -83,34 +87,116 @@ func (r *Router) handleAssignmentBind(ctx context.Context, env proto.Envelope) e var input proto.AssignmentBindPayload ref, code := env.Assignment, "" if env.ID == "" || env.DecodeRequest(&input) != nil || input.Validate() != nil || !ref.Valid() { - code = "invalid_request" - } else { - var resource sandboxbootstrap.Resource - if input.Resource != nil { - resource = *input.Resource - } - r.mu.Lock() - a := r.assignments[ref.SessionID] - switch { - case a == nil: - var environment Environment - if r.environments != nil { - if environment = r.environments(ref, input); environment == nil { - code = proto.AssignmentConflict - break - } + return r.reply(ctx, env, proto.TypeAssignmentStatus, assignmentStatus(proto.AssignmentBound, "invalid_request")) + } + var resource sandboxbootstrap.Resource + if input.Resource != nil { + resource = *input.Resource + } + r.mu.Lock() + if r.suspensions[input.EnvironmentID] != nil { + r.mu.Unlock() + return ErrRouterQuiesced + } + a := r.assignments[ref.SessionID] + switch { + case a == nil: + 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): + } + r.assignments[ref.SessionID] = &assignmentState{ref: ref, environmentID: input.EnvironmentID, resource: resource, grant: input.AttachGrant, environment: environment} + case a.ref.AssignmentID != ref.AssignmentID: + code = proto.AssignmentConflict + case ref.Epoch > a.ref.Epoch && (!a.released || a.superseding != nil): + // The bind supersedes the earlier epoch, or a pending supersede. + preparations := r.fenceSessionWorkLocked(ref.SessionID) + a.ref, a.released, a.superseding = ref, true, &input + r.shutdownWG.Add(1) + r.mu.Unlock() + go r.supersede(ctx, env, a, preparations) + return nil + case ref.Epoch == a.ref.Epoch && a.superseding != nil: + if !sameBind(*a.superseding, input) { code = proto.AssignmentConflict + break } + // A retry repeats the pending supersede. + r.shutdownWG.Add(1) r.mu.Unlock() + go r.supersede(ctx, env, a, nil) + return nil + case 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): + code = proto.AssignmentConflict } + r.mu.Unlock() return r.reply(ctx, env, proto.TypeAssignmentStatus, assignmentStatus(proto.AssignmentBound, code)) } +// supersede settles the earlier epoch's work, closes the Session's Executor +// and owner, and then binds the pending supersede at env's assignment. It +// replies assignment_stale when a release or a later bind took over meanwhile. +// A failed cleanup keeps the supersede pending, so a retry repeats it. +func (r *Router) supersede(ctx context.Context, env proto.Envelope, a *assignmentState, preparations []*preparationState) { + defer r.shutdownWG.Done() + ref := env.Assignment + cleanupCtx, stop := r.shutdownContext(context.WithoutCancel(ctx)) + defer stop() + for _, p := range preparations { + r.releasePreparation(p, "failed", proto.AssignmentStale, true) + } + a.work.Wait() + a.cleanup.Lock() + defer a.cleanup.Unlock() + r.mu.Lock() + pending, environment := a.ref == ref && a.superseding != nil, a.environment + r.mu.Unlock() + var err error + if pending { + err = r.closeSessionExecutor(ref.SessionID) + } + if pending && err == nil && environment != nil { + err = environment.Close(cleanupCtx) + } + code := "" + r.mu.Lock() + switch { + case a.ref != ref || a.released && a.superseding == nil: + code = proto.AssignmentStale + case a.superseding == nil: + // An earlier attempt bound it. + case err != nil: + r.log.Warn("assignment supersede cleanup unconfirmed", "session_id", ref.SessionID, "err", err) + code = proto.CleanupUnconfirmed + default: + input := *a.superseding + var owner Environment + if r.environments != nil { + if owner = r.environments(ref, input); owner == nil { + code = proto.AssignmentConflict + break + } + } + var resource sandboxbootstrap.Resource + if input.Resource != nil { + resource = *input.Resource + } + a.environmentID, a.resource, a.grant, a.environment = input.EnvironmentID, resource, input.AttachGrant, owner + a.released, a.superseding = false, nil + } + r.mu.Unlock() + _ = r.reply(cleanupCtx, env, proto.TypeAssignmentStatus, assignmentStatus(proto.AssignmentBound, code)) +} + +func sameBind(a, b proto.AssignmentBindPayload) bool { + return a.EnvironmentID == b.EnvironmentID && (a.Resource == nil) == (b.Resource == nil) && (a.Resource == nil || *a.Resource == *b.Resource) && bytes.Equal(a.AttachGrant, b.AttachGrant) +} + // handleAssignmentRelease fences the assignment, then settles the Session's // work and Executor, closes its Environment owner and removes its home before // it replies. A retry at the same epoch repeats the cleanup. @@ -135,7 +221,7 @@ func (r *Router) handleAssignmentRelease(ctx context.Context, env proto.Envelope case ref.Epoch < a.ref.Epoch: code = proto.AssignmentStale default: - a.ref, a.released = ref, true + a.ref, a.released, a.superseding = ref, true, nil } if code != "" { r.mu.Unlock() @@ -179,11 +265,11 @@ func (r *Router) handleAssignmentRelease(ctx context.Context, env proto.Envelope // the release releases, which also cancels their exports. Router.mu must be // held. func (r *Router) fenceSessionWorkLocked(sessionID string) []*preparationState { - if u := r.workspaceWrite; u != nil && u.envelope.Assignment.SessionID == sessionID && !u.finished { + if u := r.workspaceWrites[sessionID]; u != nil && !u.finished { u.finished = true close(u.ready) } - if u := r.runtimePreparation; u != nil && u.envelope.Assignment.SessionID == sessionID && !u.finished { + if u := r.runtimePreparations[sessionID]; u != nil && !u.finished { r.finishRuntimePreparationTransferLocked(u, false) } var preparations []*preparationState diff --git a/apps/daemon/internal/dispatch/assignment_test.go b/apps/daemon/internal/dispatch/assignment_test.go index 930348dae..631d76923 100644 --- a/apps/daemon/internal/dispatch/assignment_test.go +++ b/apps/daemon/internal/dispatch/assignment_test.go @@ -262,3 +262,60 @@ func TestReleaseRetryAndShutdownCloseOwnersOnce(t *testing.T) { t.Fatalf("released owner closes = %d (overlapped %t), unreleased owner closes = %d", released.closes.Load(), released.overlapped.Load(), unreleased.closes.Load()) } } + +func TestSupersedingBindFencesTheEarlierEpoch(t *testing.T) { + h := newHarness(t) + defer h.router.Shutdown(context.Background()) + startRun(t, h.router, h.sender, "fake_alpha", "s") + sess := <-h.gotSess + bind := func(id string, epoch uint64, payload proto.AssignmentBindPayload) proto.AssignmentStatusPayload { + env := scoped(t, "s", proto.TypeAssignmentBind, id, payload) + env.Assignment.Epoch = epoch + if err := h.router.Handle(t.Context(), env); err != nil { + t.Fatal(err) + } + return waitAssignmentStatus(t, h.sender, id) + } + if got := bind("supersede", 2, proto.AssignmentBindPayload{}); got.State != proto.AssignmentBound || got.ErrorCode != "" { + t.Fatalf("superseding bind = %+v", got) + } + // The earlier epoch's Run ended before the bind replied. + terminal, bound := -1, -1 + for i, frame := range h.sender.snapshot() { + switch { + case frame.ID == "s" && (frame.Type == proto.TypeDone || frame.Type == proto.TypeError): + terminal = i + case frame.ID == "supersede": + bound = i + } + } + if terminal < 0 || terminal > bound || sess.cancels() != 1 { + t.Fatalf("Run terminal at %d, bind reply at %d, cancels = %d", terminal, bound, sess.cancels()) + } + for id, test := range map[string]struct { + epoch uint64 + payload proto.AssignmentBindPayload + code string + }{ + "lower": {1, proto.AssignmentBindPayload{}, proto.AssignmentStale}, + "changed": {2, proto.AssignmentBindPayload{EnvironmentID: uuid.NewString()}, proto.AssignmentConflict}, + "repeated": {2, proto.AssignmentBindPayload{}, ""}, + } { + if got := bind(id, test.epoch, test.payload); got.ErrorCode != test.code { + t.Fatalf("%s bind = %+v", id, got) + } + } + stale := scoped(t, "s", proto.TypeExecutionPrepare, "stale", noEnvironmentPreparation("s", proto.PromptRequestPayload{AgentKind: "fake_alpha"})) + if err := h.router.Handle(t.Context(), stale); err == nil { + t.Fatal("the superseded epoch admitted a preparation") + } + if got := waitPreparationStatus(t, h.sender, "stale", "rejected", ""); got.ErrorCode != proto.AssignmentStale { + t.Fatalf("superseded preparation = %+v", got) + } + current := stale + current.ID, current.Assignment.Epoch = "current", 2 + if err := h.router.Handle(t.Context(), current); err != nil { + t.Fatal(err) + } + waitPreparationStatus(t, h.sender, "current", "ready", "") +} diff --git a/apps/daemon/internal/dispatch/cancellation.go b/apps/daemon/internal/dispatch/cancellation.go index c54e9eb45..adcc02015 100644 --- a/apps/daemon/internal/dispatch/cancellation.go +++ b/apps/daemon/internal/dispatch/cancellation.go @@ -30,8 +30,12 @@ func (r *Router) handlePromptCancel(ctx context.Context, env proto.Envelope) err handoff := state.preparedHandoff release, attempt := r.claimPreparedReleaseLocked(state, true, "", true) if request.DeliveryID != "" { + done := r.trackWorkLocked(env.Assignment) r.shutdownWG.Add(1) - go r.sendPreparedCancellation(state, handoff, release, attempt, env, request.DeliveryID) + go func() { + defer done() + r.sendPreparedCancellation(state, handoff, release, attempt, env, request.DeliveryID) + }() } r.mu.Unlock() return nil diff --git a/apps/daemon/internal/dispatch/executor.go b/apps/daemon/internal/dispatch/executor.go index 72b7c7e56..bc75c1f69 100644 --- a/apps/daemon/internal/dispatch/executor.go +++ b/apps/daemon/internal/dispatch/executor.go @@ -83,7 +83,7 @@ func (r *Router) handleExecutorPrepare(ctx context.Context, env proto.Envelope, requestFingerprint := sha256.Sum256(encoded) req.Assignment = env.Assignment r.mu.Lock() - if r.closed || r.suspension != nil { + if r.closed { r.mu.Unlock() return ErrRouterClosed } @@ -208,7 +208,7 @@ func (r *Router) prepareExecutor(p *preparationState, req proto.PromptRequestPay owner.native, owner.preparing = native, false close(owner.prepared) r.log.Info("executor native_prepare", "executor_id", owner.id, "session_id", owner.sessionID, "duration_ms", time.Since(started).Milliseconds(), "success", err == nil && native != nil) - ready := native != nil && err == nil && !owner.invalid && !r.closed && r.suspension == nil && p.status.State == "preparing" && p.ctx.Err() == nil + ready := native != nil && err == nil && !owner.invalid && !r.closed && p.status.State == "preparing" && p.ctx.Err() == nil p.busy = false if ready { p.status.State, p.status.Revision = "ready", p.status.Revision+1 diff --git a/apps/daemon/internal/dispatch/native_file_results_test.go b/apps/daemon/internal/dispatch/native_file_results_test.go index a3d15e8bd..4f5bc685a 100644 --- a/apps/daemon/internal/dispatch/native_file_results_test.go +++ b/apps/daemon/internal/dispatch/native_file_results_test.go @@ -37,9 +37,10 @@ func TestLocalUploadUnknownRetainsOwner(t *testing.T) { if got.Outcome != "unknown" { t.Fatal(got) } - r.workspaceWrite = &workspaceUpload{envelope: proto.Envelope{ID: uuid.NewString()}, finished: true, uncertain: true} request := proto.WorkspaceWritePayload{Step: "begin", EnvironmentID: uuid.NewString(), SessionID: uuid.NewString(), Path: "file", SizeBytes: 0, SHA256: "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"} + r.workspaceWrites[request.SessionID] = &workspaceUpload{envelope: proto.Envelope{ID: uuid.NewString()}, finished: true, uncertain: true} env, _ := proto.NewEnvelope(proto.TypeWorkspaceWrite, uuid.NewString(), request) + env.Assignment.SessionID = request.SessionID if err = r.Handle(t.Context(), env); err != nil { t.Fatal(err) } diff --git a/apps/daemon/internal/dispatch/preparation.go b/apps/daemon/internal/dispatch/preparation.go index 5f31485a5..76ca2ab2b 100644 --- a/apps/daemon/internal/dispatch/preparation.go +++ b/apps/daemon/internal/dispatch/preparation.go @@ -83,7 +83,7 @@ func (r *Router) handleExecutionPrepare(ctx context.Context, env proto.Envelope) r.mu.Unlock() return r.rejectPreparation(env, code) } - if u := r.runtimePreparation; u != nil && u.envelope.Assignment.SessionID == input.SessionID { + if r.runtimePreparations[input.SessionID] != nil { r.mu.Unlock() return r.rejectPreparation(env, "resource_unavailable") } @@ -168,7 +168,7 @@ func (r *Router) handleExecutionRelease(_ context.Context, env proto.Envelope) e func (r *Router) releasePreparation(p *preparationState, state, code string, publish bool) { r.mu.Lock() - if r.closed || r.suspension != nil { + if r.closed { r.mu.Unlock() return } @@ -221,7 +221,7 @@ func (r *Router) publishPreparation(p *preparationState, status proto.Preparatio r.mu.Unlock() return } - if r.closed || r.suspension != nil { + if r.closed { r.mu.Unlock() return } diff --git a/apps/daemon/internal/dispatch/preparation_start.go b/apps/daemon/internal/dispatch/preparation_start.go index 557234eca..c1e268d98 100644 --- a/apps/daemon/internal/dispatch/preparation_start.go +++ b/apps/daemon/internal/dispatch/preparation_start.go @@ -20,11 +20,11 @@ func (r *Router) handleExecutionStart(_ context.Context, env proto.Envelope) err encoded, _ := json.Marshal(input) fingerprint := sha256.Sum256(encoded) r.mu.Lock() - if r.closed || r.suspension != nil { + if r.closed { r.mu.Unlock() return ErrRouterClosed } - if r.workspaceWrite != nil || r.workspaceExport != nil || r.runtimePreparation != nil { + if r.environmentTransferLocked(env.Assignment.SessionID) { r.mu.Unlock() return r.rejectPreparation(env, "resource_unavailable") } diff --git a/apps/daemon/internal/dispatch/prepared_handoff.go b/apps/daemon/internal/dispatch/prepared_handoff.go index 772cce4a5..e5fdf3856 100644 --- a/apps/daemon/internal/dispatch/prepared_handoff.go +++ b/apps/daemon/internal/dispatch/prepared_handoff.go @@ -236,6 +236,8 @@ func (r *Router) runPreparedRelease(state *sessionState, handoff *preparedHandof r.mu.Unlock() r.mu.Lock() outputErr, terminal, closed := handoff.outputErr, handoff.terminal, r.closed + // The retired Run's terminal is the Session's work until it is sent. + tail := r.trackWorkLocked(state.assignment) r.mu.Unlock() // Done can become visible before Sender.Send returns. Commit the settled // owner and retire this Run before publication so an immediate successor @@ -281,6 +283,7 @@ func (r *Router) runPreparedRelease(state *sessionState, handoff *preparedHandof p.closeErr = terminalErr close(release.settled) r.mu.Unlock() + tail() } func (r *Router) forwardPreparedOutput(state *sessionState) { diff --git a/apps/daemon/internal/dispatch/router.go b/apps/daemon/internal/dispatch/router.go index 138ffe3be..1fe1c9ef2 100644 --- a/apps/daemon/internal/dispatch/router.go +++ b/apps/daemon/internal/dispatch/router.go @@ -33,9 +33,8 @@ type Router struct { log *slog.Logger admission sync.RWMutex - suspension *proto.EnvironmentSuspendPayload - suspendedBy proto.AssignmentRef // the assignment that quiesced mu sync.Mutex + suspensions map[string]*suspension // Environment ID → its quiescence assignments map[string]*assignmentState // SessionID → assignment sessions map[string]*sessionState // RunID → state applied map[string]appliedFunctionResult // RunID and call ID → applied result @@ -48,9 +47,12 @@ type Router struct { preparations map[string]*preparationState preparationRequests map[string]*preparationState preparationTimeout time.Duration - runtimePreparation *runtimePreparationTransfer - workspaceWrite *workspaceUpload - workspaceExport *workspaceExport + // Each Session has at most one transfer of each kind; transferBytes + // counts the bodies that admitted writes and Runtime preparations buffer. + runtimePreparations map[string]*runtimePreparationTransfer // SessionID → + workspaceWrites map[string]*workspaceUpload // SessionID → + workspaceExports map[string]*workspaceExport // SessionID → + transferBytes int workspaceReads map[string]string // read ID → SessionID environments func(proto.AssignmentRef, proto.AssignmentBindPayload) Environment removeHome func(sessionID string) error @@ -91,9 +93,10 @@ type Config struct { IdleTimeout time.Duration PreparationTimeout time.Duration // Environments resolves the Environment owner of a Session's first bind on - // this Router, under the Router's lock, without I/O. A nil owner rejects - // the bind: the Runtime does not serve that Session. Nil Environments - // leaves every Session without an owner. + // this Router, and of a bind that supersedes its assignment, under the + // Router's lock, without I/O. 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 // RemoveHome removes the Session's native home once its Executors have // closed. Nil declares that assignment_release does not accept RemoveHome. @@ -136,19 +139,34 @@ func New(cfg Config) (*Router, error) { preparations: make(map[string]*preparationState), preparationRequests: make(map[string]*preparationState), preparationTimeout: preparationTimeout, + suspensions: make(map[string]*suspension), + runtimePreparations: make(map[string]*runtimePreparationTransfer), + workspaceWrites: make(map[string]*workspaceUpload), + workspaceExports: make(map[string]*workspaceExport), environments: cfg.Environments, removeHome: cfg.RemoveHome, }, nil } +// transferMemory bounds the bodies that a connection's admitted workspace +// writes and Runtime preparations buffer together. +const transferMemory = proto.WorkspaceWriteMaxBytes + proto.RuntimePrepareMaxBytes + // Handle dispatches one inbound Envelope. Errors are returned for // programmer-visible problems (bad shape, registry miss); transient -// session-level failures are logged and swallowed. A quiesced Router -// still handles assignment_release. +// session-level failures are logged and swallowed. A quiesced Environment's +// Sessions still have their assignment_release handled. // // Adopts env.Trace into ctx so every downstream log under it inherits // the same trace_id, making a single grep cover both sides. func (r *Router) Handle(ctx context.Context, env proto.Envelope) error { + ctx = adoptEnvelopeTrace(ctx, env) + switch env.Type { + case proto.TypeEnvironmentQuiesce: + return r.handleQuiesce(ctx, env) + case proto.TypeEnvironmentResume: + return r.handleResume(ctx, env) + } r.admission.RLock() defer r.admission.RUnlock() r.mu.Lock() @@ -156,14 +174,12 @@ func (r *Router) Handle(ctx context.Context, env proto.Envelope) error { r.mu.Unlock() return ErrRouterClosed } - if r.suspension != nil && env.Type != proto.TypeAssignmentRelease { + if a := r.assignments[env.Assignment.SessionID]; a != nil && r.suspensions[a.environmentID] != nil && env.Type != proto.TypeAssignmentRelease { r.mu.Unlock() return ErrRouterQuiesced } r.mu.Unlock() - ctx = adoptEnvelopeTrace(ctx, env) - switch env.Type { case proto.TypeAssignmentBind: return r.handleAssignmentBind(ctx, env) diff --git a/apps/daemon/internal/dispatch/runtime_preparation.go b/apps/daemon/internal/dispatch/runtime_preparation.go index 8db2ec1f4..f8ca616dd 100644 --- a/apps/daemon/internal/dispatch/runtime_preparation.go +++ b/apps/daemon/internal/dispatch/runtime_preparation.go @@ -14,8 +14,9 @@ import ( const runtimePreparationTimeout = 120 * time.Second -// Router.mu protects one connection-local transfer. Partial installation data -// belongs to the bound Environment and is never removed by transfer cleanup. +// Router.mu protects a Session's Runtime preparation transfer. Partial +// installation data belongs to the bound Environment and is never removed by +// transfer cleanup. type runtimePreparationTransfer struct { id uuid.UUID // the envelope's ID envelope proto.Envelope @@ -36,9 +37,10 @@ func (r *Router) handleRuntimePrepare(ctx context.Context, env proto.Envelope) e var request proto.RuntimePreparePayload if len(env.Payload) > proto.RuntimePrepareMaxFrameBytes || env.DecodeRequest(&request) != nil || !proto.ValidRuntimePrepareRequest(request) { r.mu.Lock() - pending := r.runtimePreparation != nil && r.runtimePreparation.envelope.ID == env.ID - if pending && !r.runtimePreparation.finished { - r.finishRuntimePreparationTransferLocked(r.runtimePreparation, false) + u := r.runtimePreparations[env.Assignment.SessionID] + pending := u != nil && u.envelope.ID == env.ID + if pending && !u.finished { + r.finishRuntimePreparationTransferLocked(u, false) } r.mu.Unlock() if pending { @@ -48,13 +50,14 @@ func (r *Router) handleRuntimePrepare(ctx context.Context, env proto.Envelope) e return r.sendRuntimePrepareResult(ctx, env, rejectedRuntimePreparation("invalid_request")) } r.mu.Lock() - if r.closed || r.suspension != nil { + if r.closed { r.mu.Unlock() return ErrRouterClosed } + u := r.runtimePreparations[env.Assignment.SessionID] if request.Step == "begin" { - if r.runtimePreparation != nil { - duplicate := r.runtimePreparation.envelope.ID == env.ID + if u != nil || r.transferBytes > transferMemory-request.SizeBytes { + duplicate := u != nil && u.envelope.ID == env.ID r.mu.Unlock() if duplicate { return errors.New("dispatch: Runtime preparation already admitted") @@ -79,7 +82,8 @@ func (r *Router) handleRuntimePrepare(ctx context.Context, env proto.Envelope) e id: id, envelope: env, request: request, data: make([]byte, 0, request.SizeBytes), ready: make(chan struct{}), cancel: cancel, } - r.runtimePreparation = u + r.runtimePreparations[request.SessionID] = u + r.transferBytes += request.SizeBytes done := r.trackWorkLocked(env.Assignment) r.shutdownWG.Add(1) r.mu.Unlock() @@ -90,7 +94,6 @@ func (r *Router) handleRuntimePrepare(ctx context.Context, env proto.Envelope) e } return nil } - u := r.runtimePreparation if u == nil || u.envelope.ID != env.ID || u.envelope.Assignment != env.Assignment { r.mu.Unlock() return r.sendRuntimePrepareResult(ctx, env, rejectedRuntimePreparation("resource_unavailable")) @@ -141,9 +144,7 @@ func (r *Router) runtimePreparationResourcesBusyLocked(sessionID string) bool { // environmentTransferLocked reports whether the Session has a workspace write, // a workspace export or a Runtime preparation. Router.mu must be held. func (r *Router) environmentTransferLocked(sessionID string) bool { - return r.workspaceWrite != nil && r.workspaceWrite.envelope.Assignment.SessionID == sessionID || - r.workspaceExport != nil && r.workspaceExport.request.Assignment.SessionID == sessionID || - r.runtimePreparation != nil && r.runtimePreparation.envelope.Assignment.SessionID == sessionID + return r.workspaceWrites[sessionID] != nil || r.workspaceExports[sessionID] != nil || r.runtimePreparations[sessionID] != nil } // sessionWorkLocked reports whether the Session has a Run, a workspace read or @@ -202,9 +203,10 @@ func (r *Router) runRuntimePreparationTransfer(ctx context.Context, u *runtimePr // Release the potentially large body before waiting on transport delivery. data = nil r.mu.Lock() + r.transferBytes -= u.request.SizeBytes u.uncertain = result.Outcome == "unknown" - if !u.uncertain && r.runtimePreparation == u { - r.runtimePreparation = nil + if !u.uncertain { + delete(r.runtimePreparations, u.request.SessionID) } r.mu.Unlock() // The result has a separate send budget, independent of an installation timeout. diff --git a/apps/daemon/internal/dispatch/runtime_preparation_test.go b/apps/daemon/internal/dispatch/runtime_preparation_test.go index da997c58f..f0a027017 100644 --- a/apps/daemon/internal/dispatch/runtime_preparation_test.go +++ b/apps/daemon/internal/dispatch/runtime_preparation_test.go @@ -103,7 +103,7 @@ func TestRuntimePreparationTransferValidatesCompleteBodyBeforeMutation(t *testin } capabilitiesReceipt(t, sender, id, "ready") r.mu.Lock() - owner := r.runtimePreparation + owner := r.runtimePreparations[session] r.mu.Unlock() duplicate := uuid.NewString() if err := r.Handle(t.Context(), capabilityEnvelope(t, duplicate, request)); err != nil { @@ -133,7 +133,7 @@ func TestRuntimePreparationTransferValidatesCompleteBodyBeforeMutation(t *testin capabilitiesReceipt(t, sender, id, "rejected") r.mu.Lock() defer r.mu.Unlock() - if owner.apply || owner.data != nil || r.runtimePreparation != nil { + if owner.apply || owner.data != nil || len(r.runtimePreparations) != 0 { t.Fatal("invalid body retained or admitted a mutation") } }) @@ -164,7 +164,7 @@ func TestRuntimePreparationBeginRequiresExactBindingAndBounds(t *testing.T) { t.Fatal(err) } capabilitiesReceipt(t, sender, id, "rejected") - if r.runtimePreparation != nil { + if len(r.runtimePreparations) != 0 { t.Fatal("invalid scope allocated a transfer") } } @@ -177,9 +177,9 @@ func TestRuntimePreparationPreparationExcludesOwnedResources(t *testing.T) { own := proto.Envelope{Assignment: capabilityRef} switch mode { case "write": - r.workspaceWrite = &workspaceUpload{envelope: own} + r.workspaceWrites[session] = &workspaceUpload{envelope: own} case "export": - r.workspaceExport = &workspaceExport{request: own} + r.workspaceExports[session] = &workspaceExport{request: own} case "read": r.workspaceReads = map[string]string{"read": session} case "run": @@ -199,17 +199,17 @@ func TestRuntimePreparationPreparationExcludesOwnedResources(t *testing.T) { if mode == "other session" { capabilitiesReceipt(t, sender, id, "ready") r.mu.Lock() - r.finishRuntimePreparationTransferLocked(r.runtimePreparation, false) + r.finishRuntimePreparationTransferLocked(r.runtimePreparations[session], false) r.mu.Unlock() capabilitiesReceipt(t, sender, id, "rejected") } else if got := capabilitiesReceipt(t, sender, id, "rejected"); got.ErrorCode != "resource_unavailable" { t.Fatal(got) } - if r.runtimePreparation != nil { + if len(r.runtimePreparations) != 0 { t.Fatal("busy Runtime admitted capability preparation") } - r.workspaceWrite = nil - r.workspaceExport = nil + clear(r.workspaceWrites) + clear(r.workspaceExports) r.workspaceReads = nil clear(r.sessions) clear(r.executors) @@ -248,10 +248,10 @@ func TestRuntimePreparationUploadBlocksWorkspaceWriteAndSuspension(t *testing.T) t.Fatal("missing write rejection") } r.mu.Lock() - owner := r.runtimePreparation + owner := r.runtimePreparations[session] r.mu.Unlock() shutdownCapabilitiesRouter(t, r) - if owner.data != nil || r.runtimePreparation != nil { + if owner.data != nil || len(r.runtimePreparations) != 0 { t.Fatal("disconnect retained uncommitted body") } } @@ -263,7 +263,7 @@ func TestRuntimePreparationCancellationKeepsOwnershipUntilApplyStops(t *testing. request := proto.RuntimePreparePayload{Step: "begin", Action: "finalize", EnvironmentID: environment, SessionID: session, Sources: &agentcapabilities.Input{}} owner := &runtimePreparationTransfer{id: uuid.MustParse(id), envelope: capabilityEnvelope(t, id, request), request: request, ready: make(chan struct{}), cancel: cancel, finished: true, apply: true} close(owner.ready) - r.runtimePreparation = owner + r.runtimePreparations[session], r.transferBytes = owner, request.SizeBytes r.shutdownWG.Add(1) started, interrupted, release := make(chan struct{}), make(chan struct{}), make(chan struct{}) retained := filepath.Join(t.TempDir(), "installed.json") @@ -286,7 +286,7 @@ func TestRuntimePreparationCancellationKeepsOwnershipUntilApplyStops(t *testing. } <-interrupted r.mu.Lock() - owned := r.runtimePreparation == owner + owned := r.runtimePreparations[session] == owner r.mu.Unlock() if !owned { t.Fatal("cancel released unsettled capability ownership") @@ -297,7 +297,7 @@ func TestRuntimePreparationCancellationKeepsOwnershipUntilApplyStops(t *testing. if _, err := os.Stat(retained); err != nil { t.Fatal("shutdown deleted installation result") } - if r.runtimePreparation != nil { + if len(r.runtimePreparations) != 0 { t.Fatal("confirmed completion retained capacity") } } @@ -344,14 +344,14 @@ func TestRuntimePreparationResultCategoriesAndUnknownOwnership(t *testing.T) { request := capabilityBegin(environment, session, []byte("abc")) owner := &runtimePreparationTransfer{id: uuid.MustParse(id), envelope: capabilityEnvelope(t, id, request), request: request, data: []byte("abc"), ready: make(chan struct{}), cancel: cancel, finished: true, apply: true} close(owner.ready) - r.runtimePreparation = owner + r.runtimePreparations[session], r.transferBytes = owner, request.SizeBytes r.shutdownWG.Add(1) go r.runRuntimePreparationTransfer(ctx, owner, func(context.Context, uuid.UUID, proto.RuntimePreparePayload, []byte) error { return context.DeadlineExceeded }, func() {}) capabilitiesReceipt(t, sender, id, "unknown") r.mu.Lock() - owned := r.runtimePreparation == owner && owner.uncertain && owner.data == nil + owned := r.runtimePreparations[session] == owner && owner.uncertain && owner.data == nil r.mu.Unlock() if !owned { t.Fatal("unknown mutation released its ownership") @@ -362,3 +362,98 @@ func TestRuntimePreparationResultCategoriesAndUnknownOwnership(t *testing.T) { t.Fatal("shutdown claimed uncertain mutation settled") } } + +func TestSessionTransferBlocksOnlyItsSession(t *testing.T) { + for _, mode := range []string{"held", "uncertain"} { + t.Run(mode, func(t *testing.T) { + r, sender, environment, session := capabilitiesTestRouter(t) + other := proto.AssignmentRef{SessionID: uuid.NewString(), AssignmentID: "other", Epoch: 1} + bindAssignment(r, other, environment) + next := func(id string) proto.Envelope { + t.Helper() + select { + case env := <-sender.frames: + if env.ID != id { + t.Fatalf("frame %s %s, want %s", env.Type, env.ID, id) + } + return env + case <-time.After(3 * time.Second): + t.Fatal("missing frame", id) + } + return proto.Envelope{} + } + if mode == "held" { + id := uuid.NewString() + if err := r.Handle(t.Context(), capabilityEnvelope(t, id, capabilityBegin(environment, session, []byte("abc")))); err != nil { + t.Fatal(err) + } + capabilitiesReceipt(t, sender, id, "ready") + } else { + r.mu.Lock() + r.workspaceWrites[session] = &workspaceUpload{finished: true, uncertain: true} + r.mu.Unlock() + } + start := func(ref proto.AssignmentRef) string { + env, err := proto.NewEnvelope(proto.TypeExecutionStart, uuid.NewString(), proto.ExecutionStartPayload{Handle: "handle", ExecutorID: "executor", RunID: "run", Input: proto.TextInput("input")}) + if err != nil { + t.Fatal(err) + } + env.Assignment = ref + _ = r.Handle(t.Context(), env) + var status proto.PreparationStatusPayload + if err := next(env.ID).DecodePayload(&status); err != nil { + t.Fatal(err) + } + return status.ErrorCode + } + if got := start(capabilityRef); got != "resource_unavailable" { + t.Fatalf("the transferring Session started a Run: %s", got) + } + if got := start(other); got != "unknown_preparation" { + t.Fatalf("another Session's transfer blocked execution_start: %s", got) + } + // The other Session's transfers proceed while the connection's + // memory bound allows their bodies. + id := uuid.NewString() + request := capabilityBegin(environment, other.SessionID, []byte("abc")) + env := capabilityEnvelope(t, id, request) + env.Assignment = other + if err := r.Handle(t.Context(), env); err != nil { + t.Fatal(err) + } + capabilitiesReceipt(t, sender, id, "ready") + r.mu.Lock() + r.finishRuntimePreparationTransferLocked(r.runtimePreparations[other.SessionID], false) + r.mu.Unlock() + capabilitiesReceipt(t, sender, id, "rejected") + digest := sha256.Sum256([]byte("abc")) + for _, buffered := range []int{transferMemory, 0} { + write, err := proto.NewEnvelope(proto.TypeWorkspaceWrite, uuid.NewString(), proto.WorkspaceWritePayload{Step: "begin", EnvironmentID: environment, SessionID: other.SessionID, Path: "proof", SizeBytes: 3, SHA256: hex.EncodeToString(digest[:])}) + if err != nil { + t.Fatal(err) + } + write.Assignment = other + r.mu.Lock() + r.transferBytes += buffered + r.mu.Unlock() + if err := r.Handle(t.Context(), write); err != nil { + t.Fatal(err) + } + r.mu.Lock() + r.transferBytes -= buffered + r.mu.Unlock() + var result proto.WorkspaceWriteResultPayload + if err := next(write.ID).DecodePayload(&result); err != nil { + t.Fatal(err) + } + if want := map[int]string{0: "ready", transferMemory: "write_capacity"}[buffered]; result.Outcome != want && result.ErrorCode != want { + t.Fatalf("write with %d bytes buffered = %+v", buffered, result) + } + } + r.mu.Lock() + delete(r.workspaceWrites, session) + r.mu.Unlock() + shutdownCapabilitiesRouter(t, r) + }) + } +} diff --git a/apps/daemon/internal/dispatch/shutdown.go b/apps/daemon/internal/dispatch/shutdown.go index 1c0e96836..d67d291f9 100644 --- a/apps/daemon/internal/dispatch/shutdown.go +++ b/apps/daemon/internal/dispatch/shutdown.go @@ -31,8 +31,8 @@ func (r *Router) Shutdown(ctx context.Context) error { victims := r.sessionCancellationsLocked() if !r.closed { r.closed = true - if r.runtimePreparation != nil { - r.runtimePreparation.cancel() + for _, u := range r.runtimePreparations { + u.cancel() } // Prepared release claims exist before this signal can interrupt output. close(r.shutdownCh) @@ -72,11 +72,15 @@ func (r *Router) runShutdownAttempt(attempt *shutdownAttempt, victims []sessionC attempt.err = errors.Join(attempt.err, err) } r.mu.Lock() - if r.runtimePreparation != nil && r.runtimePreparation.uncertain { - attempt.err = errors.Join(attempt.err, errors.New("dispatch: capability preparation remains uncertain")) + for session, u := range r.runtimePreparations { + if u.uncertain { + attempt.err = errors.Join(attempt.err, fmt.Errorf("dispatch: capability preparation of Session %s remains uncertain", session)) + } } - if r.workspaceWrite != nil && r.workspaceWrite.uncertain { - attempt.err = errors.Join(attempt.err, errors.New("dispatch: local workspace write remains uncertain")) + for session, u := range r.workspaceWrites { + if u.uncertain { + attempt.err = errors.Join(attempt.err, fmt.Errorf("dispatch: local workspace write of Session %s remains uncertain", session)) + } } for _, p := range r.preparations { if p.owns { diff --git a/apps/daemon/internal/dispatch/suspend.go b/apps/daemon/internal/dispatch/suspend.go index 6ff583c86..ed9fed0e0 100644 --- a/apps/daemon/internal/dispatch/suspend.go +++ b/apps/daemon/internal/dispatch/suspend.go @@ -4,6 +4,7 @@ import ( "context" "errors" "strings" + "time" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) @@ -17,11 +18,116 @@ type AssignmentError string func (e AssignmentError) Error() string { return "dispatch: " + string(e) } -// Quiesce serializes against admission, then drains every admitted output and -// receipt and the Environment's owners before acknowledging suspension. Busy rejection leaves admission open; -// a drain timeout keeps it closed until the caller shuts the connection down. -// ref must admit work in the suspended Environment. +// quiesceTimeout bounds the drain of an environment_quiesce that Handle +// admitted. +const quiesceTimeout = 5 * time.Second + +// suspension is one quiesced Environment: the request and the assignment +// that quiesced it, which its resume must carry. +type suspension struct { + request proto.EnvironmentSuspendPayload + by proto.AssignmentRef +} + +// Quiesce quiesces the request's Environment for a caller that then +// suspends the whole Runtime: after fencing the Environment it waits for every +// output and receipt admitted on the connection, then closes the +// Environment's Executors and owners before it returns. Busy rejection leaves +// admission open; a failed drain keeps the Environment quiesced until a +// matching resume or shutdown. ref must admit work in the Environment. func (r *Router) Quiesce(ctx context.Context, ref proto.AssignmentRef, request proto.EnvironmentSuspendPayload) error { + if err := r.fenceEnvironment(ref, request); err != nil { + return err + } + if err := r.shutdownWG.waitContext(ctx); err != nil { + return err + } + return r.drainEnvironment(ctx, request.EnvironmentID) +} + +// Resume reopens the quiesced Environment and replaces the sender, after the +// caller authenticated a new connection and Core confirmed the exact +// suspension on it under the assignment that quiesced. +func (r *Router) Resume(ref proto.AssignmentRef, request proto.EnvironmentSuspendPayload, sender Sender) error { + if sender == nil { + return errors.New("dispatch: no sender") + } + r.mu.Lock() + defer r.mu.Unlock() + if err := r.resumeLocked(ref, request); err != nil { + return err + } + r.sender = sender + return nil +} + +// handleQuiesce quiesces an Environment for Core on this connection. It +// replies environment_quiesced once the drain settles, without blocking +// Handle; a failed drain replies resource_busy. +func (r *Router) handleQuiesce(ctx context.Context, env proto.Envelope) error { + var request proto.EnvironmentSuspendPayload + if env.ID == "" || len(env.ID) > 128 || env.DecodeRequest(&request) != nil || request.Rollback { + return r.replySuspension(ctx, env, proto.TypeEnvironmentQuiesced, request, "invalid_request") + } + if err := r.fenceEnvironment(env.Assignment, request); err != nil { + return r.replySuspension(ctx, env, proto.TypeEnvironmentQuiesced, request, suspensionCode(err)) + } + r.shutdownWG.Add(1) + go func() { + defer r.shutdownWG.Done() + drain, stop := r.shutdownContext(context.WithoutCancel(ctx)) + defer stop() + drain, cancel := context.WithTimeout(drain, quiesceTimeout) + defer cancel() + code := "" + if err := r.drainEnvironment(drain, request.EnvironmentID); err != nil { + r.log.Warn("environment quiesce unconfirmed", "environment_id", request.EnvironmentID, "err", err) + code = "resource_busy" + } + _ = r.replySuspension(context.WithoutCancel(ctx), env, proto.TypeEnvironmentQuiesced, request, code) + }() + return nil +} + +// handleResume reopens the quiesced Environment that the resume names. A +// rollback of an Environment this connection has not quiesced is accepted +// and changes nothing. +func (r *Router) handleResume(ctx context.Context, env proto.Envelope) error { + var request proto.EnvironmentSuspendPayload + if env.ID == "" || len(env.ID) > 128 || env.DecodeRequest(&request) != nil { + return r.replySuspension(ctx, env, proto.TypeEnvironmentResumed, request, "invalid_request") + } + r.mu.Lock() + err := r.resumeLocked(env.Assignment, request) + if errors.Is(err, errNotSuspended) && request.Rollback { + err = nil + } + r.mu.Unlock() + return r.replySuspension(ctx, env, proto.TypeEnvironmentResumed, request, suspensionCode(err)) +} + +var errNotSuspended = errors.New("dispatch: environment not suspended") + +// resumeLocked reopens the Environment that request quiesced under ref. +// Router.mu must be held. +func (r *Router) resumeLocked(ref proto.AssignmentRef, request proto.EnvironmentSuspendPayload) error { + s := r.suspensions[request.EnvironmentID] + switch { + case r.closed: + return ErrRouterClosed + case s == nil || !s.request.SameSuspension(request): + return errNotSuspended + case ref != s.by: + return AssignmentError(proto.AssignmentConflict) + } + delete(r.suspensions, request.EnvironmentID) + return nil +} + +// fenceEnvironment quiesces the request's Environment once none of its +// Sessions has unsettled work: from then on Handle admits only their +// releases and the matching resume. +func (r *Router) fenceEnvironment(ref proto.AssignmentRef, request proto.EnvironmentSuspendPayload) error { if strings.TrimSpace(request.EnvironmentID) == "" || strings.TrimSpace(request.SuspendID) == "" || len(request.SuspendID) > 128 { return errors.New("dispatch: invalid suspension identity") } @@ -30,80 +136,103 @@ func (r *Router) Quiesce(ctx context.Context, ref proto.AssignmentRef, request p } defer r.admission.Unlock() r.mu.Lock() - if r.closed { - r.mu.Unlock() + defer r.mu.Unlock() + switch { + case r.closed: return ErrRouterClosed - } - if r.suspension != nil { - r.mu.Unlock() + case r.suspensions[request.EnvironmentID] != nil: return ErrRouterQuiesced } if code := r.admitLocked(ref, ref.SessionID, request.EnvironmentID); code != "" { - r.mu.Unlock() return AssignmentError(code) } - if r.runtimePreparation != nil || len(r.sessions) != 0 || len(r.workspaceReads) != 0 || r.workspaceWrite != nil || r.workspaceExport != nil { - r.mu.Unlock() - return ErrRouterBusy + in := func(sessionID string) bool { + a := r.assignments[sessionID] + return a != nil && a.environmentID == request.EnvironmentID + } + for session := range r.assignments { + if in(session) && r.environmentTransferLocked(session) { + return ErrRouterBusy + } + } + for _, state := range r.sessions { + if state.environmentID == request.EnvironmentID { + return ErrRouterBusy + } + } + for _, session := range r.workspaceReads { + if in(session) { + return ErrRouterBusy + } } for _, p := range r.preparations { - if p.owns || p.busy { - r.mu.Unlock() + if (p.owns || p.busy) && p.environmentID == request.EnvironmentID { return ErrRouterBusy } } for _, owner := range r.executors { - if owner.invalid || owner.preparing || owner.run != nil || owner.admission != nil || owner.environmentID != request.EnvironmentID { - r.mu.Unlock() + if owner.environmentID == request.EnvironmentID && (owner.invalid || owner.preparing || owner.run != nil || owner.admission != nil) { return ErrRouterBusy } } - r.suspension, r.suspendedBy = &request, ref - for _, p := range r.preparations { - if p.timer != nil { - p.timer.Stop() + r.suspensions[request.EnvironmentID] = &suspension{request: request, by: ref} + return nil +} + +// drainEnvironment waits until the quiesced Environment's Sessions have sent +// every result they owe, then closes their Executors and releases their +// Environment owners. +func (r *Router) drainEnvironment(ctx context.Context, environmentID string) error { + r.mu.Lock() + var work []*dispatchWork + for _, a := range r.assignments { + if a.environmentID == environmentID { + work = append(work, &a.work) } } + var executors []*executorState for _, owner := range r.executors { - owner.idleLease++ - if owner.timer != nil { - owner.timer.Stop() + if owner.environmentID == environmentID { + owner.invalid, owner.closeReason = true, "quiesced" + executors = append(executors, owner) } } r.mu.Unlock() - err := r.shutdownWG.waitContext(ctx) - if err == nil { - err = r.closeEnvironments(ctx, request.EnvironmentID) + for _, w := range work { + if err := w.waitContext(ctx); err != nil { + return err + } } - r.mu.Lock() - if r.closed { - err = ErrRouterClosed + for _, owner := range executors { + if err := r.closeExecutor(owner); err != nil { + return err + } + } + if err := r.closeEnvironments(ctx, environmentID); err != nil { + return err } - r.mu.Unlock() - return err -} - -// Resume opens admission only after the caller authenticated a new connection -// and Core confirmed the exact suspension identity on that connection under the -// assignment that quiesced. -func (r *Router) Resume(ref proto.AssignmentRef, request proto.EnvironmentSuspendPayload, sender Sender) error { - r.admission.Lock() - defer r.admission.Unlock() r.mu.Lock() defer r.mu.Unlock() if r.closed { return ErrRouterClosed } - if sender == nil || r.suspension == nil || !r.suspension.SameSuspension(request) { - return errors.New("dispatch: suspension identity mismatch") - } - if ref != r.suspendedBy { - return AssignmentError(proto.AssignmentConflict) - } - r.sender = sender - r.suspension, r.suspendedBy = nil, proto.AssignmentRef{} - for _, owner := range r.executors { - r.scheduleExecutorIdleLocked(owner) - } return nil } + +// suspensionCode is the error code of a refused quiesce or resume. +func suspensionCode(err error) string { + var rejected AssignmentError + switch { + case err == nil: + return "" + case errors.As(err, &rejected): + return string(rejected) + case errors.Is(err, errNotSuspended): + return "not_suspended" + } + return "resource_busy" +} + +func (r *Router) replySuspension(ctx context.Context, env proto.Envelope, typ string, request proto.EnvironmentSuspendPayload, code string) error { + return r.reply(ctx, env, typ, proto.EnvironmentSuspendResultPayload{EnvironmentID: request.EnvironmentID, SuspendID: request.SuspendID, Accepted: code == "", ErrorCode: code}) +} diff --git a/apps/daemon/internal/dispatch/suspend_test.go b/apps/daemon/internal/dispatch/suspend_test.go index efe1d5714..482c6fe93 100644 --- a/apps/daemon/internal/dispatch/suspend_test.go +++ b/apps/daemon/internal/dispatch/suspend_test.go @@ -52,13 +52,16 @@ func suspensionRouter(t *testing.T, sender Sender) *Router { } func TestQuiesceRejectsEveryUnsettledResource(t *testing.T) { + session := suspendRef.SessionID cases := map[string]func(*Router){ - "active": func(r *Router) { r.sessions["run"] = &sessionState{} }, - "preparing": func(r *Router) { r.preparations["p"] = &preparationState{owns: true} }, - "receipt": func(r *Router) { r.preparations["p"] = &preparationState{busy: true} }, - "read": func(r *Router) { r.workspaceReads = map[string]string{"read": suspendRef.SessionID} }, - "write": func(r *Router) { r.workspaceWrite = &workspaceUpload{} }, - "export": func(r *Router) { r.workspaceExport = &workspaceExport{} }, + "active": func(r *Router) { r.sessions["run"] = &sessionState{environmentID: "env"} }, + "preparing": func(r *Router) { r.preparations["p"] = &preparationState{owns: true, environmentID: "env"} }, + "receipt": func(r *Router) { r.preparations["p"] = &preparationState{busy: true, environmentID: "env"} }, + "read": func(r *Router) { r.workspaceReads = map[string]string{"read": session} }, + "write": func(r *Router) { r.workspaceWrites[session] = &workspaceUpload{} }, + "export": func(r *Router) { r.workspaceExports[session] = &workspaceExport{} }, + "runtime preparation": func(r *Router) { r.runtimePreparations[session] = &runtimePreparationTransfer{} }, + "executor": func(r *Router) { r.executors[session] = &executorState{environmentID: "env", preparing: true} }, } for name, setup := range cases { t.Run(name, func(t *testing.T) { @@ -72,8 +75,10 @@ func TestQuiesceRejectsEveryUnsettledResource(t *testing.T) { r.sessions = map[string]*sessionState{} r.preparations = map[string]*preparationState{} r.workspaceReads = nil - r.workspaceWrite = nil - r.workspaceExport = nil + clear(r.workspaceWrites) + clear(r.workspaceExports) + clear(r.runtimePreparations) + clear(r.executors) }) } } @@ -90,7 +95,7 @@ func TestQuiesceDrainsPendingReceiptAndFencesConcurrentAdmission(t *testing.T) { deadline := time.After(time.Second) for { r.mu.Lock() - parked := r.suspension != nil + parked := r.suspensions["env"] != nil r.mu.Unlock() if parked { break @@ -151,33 +156,83 @@ func TestResumeRequiresExactSuspensionAndAssignment(t *testing.T) { } } -func TestQuiescePreservesIdleExecutorAgainstExpiredTimer(t *testing.T) { - sender := suspendSender(func(context.Context, proto.Envelope) error { return nil }) - r := suspensionRouter(t, sender) - native := &suspendedExecutor{} - owner := &executorState{id: "executor", sessionID: "session", environmentID: "env", native: native, cancel: func() {}} +func TestQuiescingOneEnvironmentLeavesAnotherRunning(t *testing.T) { + frames := make(chan proto.Envelope, 16) + r := suspensionRouter(t, suspendSender(func(_ context.Context, env proto.Envelope) error { frames <- env; return nil })) + other := proto.AssignmentRef{SessionID: "other", AssignmentID: "other", Epoch: 1} + bindAssignment(r, other, "other") + quiesced, running := &suspendedExecutor{}, &suspendedExecutor{} r.mu.Lock() - r.executors[owner.sessionID] = owner - r.scheduleExecutorIdleLocked(owner) - oldLease := owner.idleLease + r.executors[suspendRef.SessionID] = &executorState{id: "quiesced", sessionID: suspendRef.SessionID, environmentID: "env", native: quiesced, cancel: func() {}} + r.executors[other.SessionID] = &executorState{id: "running", sessionID: other.SessionID, environmentID: "other", native: running, cancel: func() {}} + // The other Environment's Session has a Run in progress. + r.sessions["run"] = &sessionState{assignment: other, environmentID: "other"} r.mu.Unlock() + suspend := func(typ, id string, ref proto.AssignmentRef, request proto.EnvironmentSuspendPayload) proto.EnvironmentSuspendResultPayload { + t.Helper() + env, err := proto.NewEnvelope(typ, id, request) + if err != nil { + t.Fatal(err) + } + env.Assignment = ref + if err := r.Handle(t.Context(), env); err != nil { + t.Fatal(err) + } + for { + select { + case frame := <-frames: + var result proto.EnvironmentSuspendResultPayload + if frame.ID == id && frame.DecodePayload(&result) == nil { + return result + } + case <-time.After(3 * time.Second): + t.Fatalf("%s %s has no result", typ, id) + } + } + } + prepare := func(ref proto.AssignmentRef) error { + return r.Handle(t.Context(), proto.Envelope{Type: proto.TypeExecutionPrepare, ID: "prepare", Assignment: ref}) + } request := proto.EnvironmentSuspendPayload{EnvironmentID: "env", SuspendID: "attempt"} - if err := r.Quiesce(context.Background(), suspendRef, request); err != nil { - t.Fatal(err) + if got := suspend(proto.TypeEnvironmentQuiesce, "quiesce", suspendRef, request); !got.Accepted { + t.Fatalf("quiesce = %+v", got) } - r.expireIdleExecutor(owner, oldLease) - if native.closed.Load() != 0 { - t.Fatal("pre-snapshot timer closed retained owner") + if quiesced.closed.Load() != 1 || running.closed.Load() != 0 { + t.Fatalf("closed Executors: quiesced %d, running %d", quiesced.closed.Load(), running.closed.Load()) } - if err := r.Resume(suspendRef, request, sender); err != nil { - t.Fatal(err) + if err := prepare(suspendRef); !errors.Is(err, ErrRouterQuiesced) { + t.Fatalf("the quiesced Environment admitted %v", err) + } + if err := prepare(other); errors.Is(err, ErrRouterQuiesced) { + t.Fatal("the other Environment stopped admitting work") + } + for id, test := range map[string]struct { + ref proto.AssignmentRef + request proto.EnvironmentSuspendPayload + code string + }{ + "foreign": {other, request, proto.AssignmentConflict}, + "unpaused": {other, proto.EnvironmentSuspendPayload{EnvironmentID: "other", SuspendID: "attempt"}, "not_suspended"}, + "rollback": {other, proto.EnvironmentSuspendPayload{EnvironmentID: "other", SuspendID: "attempt", Rollback: true}, ""}, + "obsolete": {suspendRef, proto.EnvironmentSuspendPayload{EnvironmentID: "env", SuspendID: "obsolete"}, "not_suspended"}, + } { + if got := suspend(proto.TypeEnvironmentResume, id, test.ref, test.request); got.ErrorCode != test.code || got.Accepted != (test.code == "") { + t.Fatalf("%s resume = %+v", id, got) + } + } + if got := suspend(proto.TypeEnvironmentResume, "resume", suspendRef, request); !got.Accepted { + t.Fatalf("resume = %+v", got) + } + if err := prepare(suspendRef); errors.Is(err, ErrRouterQuiesced) { + t.Fatal("resume did not reopen the Environment") } r.mu.Lock() - newLease := owner.idleLease + _, run := r.sessions["run"] + // The synthetic Run owns no goroutine; remove it before cleanup. + clear(r.sessions) r.mu.Unlock() - r.expireIdleExecutor(owner, newLease) - if native.closed.Load() != 1 { - t.Fatal("normal idle expiration was not restored") + if !run || running.closed.Load() != 0 { + t.Fatal("quiescing one Environment ended another's work") } } diff --git a/apps/daemon/internal/dispatch/workspace_export.go b/apps/daemon/internal/dispatch/workspace_export.go index eeb2d912c..e7271f30e 100644 --- a/apps/daemon/internal/dispatch/workspace_export.go +++ b/apps/daemon/internal/dispatch/workspace_export.go @@ -26,7 +26,7 @@ func (r *Router) handleWorkspaceExport(ctx context.Context, env proto.Envelope) r.mu.Unlock() return ErrRouterClosed } - u := r.workspaceExport + u := r.workspaceExports[env.Assignment.SessionID] if request.Step != "begin" { if u == nil || u.request.ID != env.ID || u.request.Assignment != env.Assignment { r.mu.Unlock() @@ -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.workspaceWrite != nil && r.workspaceWrite.envelope.Assignment.SessionID == env.Assignment.SessionID || p.executor != nil) { + if code == "" && (u != nil || r.workspaceWrites[env.Assignment.SessionID] != nil || p.executor != nil) { code = "resource_unavailable" } if code != "" { @@ -60,7 +60,7 @@ func (r *Router) handleWorkspaceExport(ctx context.Context, env proto.Envelope) owner, cancel := context.WithTimeout(owner, 180*time.Second) u = &workspaceExport{request: env, requests: make(chan proto.WorkspaceExportPayload, 1), cancel: func() { cancel(); stop() }} u.requests <- request - r.workspaceExport = u + r.workspaceExports[env.Assignment.SessionID] = u done := r.trackWorkLocked(env.Assignment) r.shutdownWG.Add(1) r.mu.Unlock() @@ -90,8 +90,8 @@ func (r *Router) runWorkspaceExport(ctx context.Context, u *workspaceExport, env _ = reader.Close() <-exported r.mu.Lock() - if r.workspaceExport == u { - r.workspaceExport = nil + if r.workspaceExports[u.request.Assignment.SessionID] == u { + delete(r.workspaceExports, u.request.Assignment.SessionID) } r.mu.Unlock() select { @@ -129,8 +129,8 @@ func (r *Router) runWorkspaceExport(ctx context.Context, u *workspaceExport, env // The next owner may start immediately after receiving completion. <-exported r.mu.Lock() - if r.workspaceExport == u { - r.workspaceExport = nil + if r.workspaceExports[u.request.Assignment.SessionID] == u { + delete(r.workspaceExports, u.request.Assignment.SessionID) } r.mu.Unlock() } diff --git a/apps/daemon/internal/dispatch/workspace_export_test.go b/apps/daemon/internal/dispatch/workspace_export_test.go index 19e9f2843..88ca1b55b 100644 --- a/apps/daemon/internal/dispatch/workspace_export_test.go +++ b/apps/daemon/internal/dispatch/workspace_export_test.go @@ -195,7 +195,7 @@ func TestWorkspaceExportCancelUnblocksWriterAndReleasesCapacity(t *testing.T) { deadline := time.Now().Add(3 * time.Second) for time.Now().Before(deadline) { r.mu.Lock() - active := r.workspaceExport != nil + active := len(r.workspaceExports) != 0 r.mu.Unlock() if !active { return @@ -220,7 +220,9 @@ func TestCanceledExportAnswersTheRequestCoreAwaits(t *testing.T) { // Core asks for the next chunk as the preparation that owns the export is released. r.mu.Lock() release() - r.workspaceExport.requests <- proto.WorkspaceExportPayload{Step: "next", Offset: int64(len(first.Data))} + for _, u := range r.workspaceExports { + u.requests <- proto.WorkspaceExportPayload{Step: "next", Offset: int64(len(first.Data))} + } r.mu.Unlock() if got := readExport(t, s); got.Outcome != "failed" && got.Outcome != "chunk" { t.Fatal("the canceled export answered with", got) diff --git a/apps/daemon/internal/dispatch/workspace_write.go b/apps/daemon/internal/dispatch/workspace_write.go index 3bf06b4e2..3b8ecefb6 100644 --- a/apps/daemon/internal/dispatch/workspace_write.go +++ b/apps/daemon/internal/dispatch/workspace_write.go @@ -11,7 +11,7 @@ import ( "github.com/google/uuid" ) -// Router.mu protects this single bounded transfer to an Environment owner. +// Router.mu protects a Session's bounded transfer to its Environment owner. type workspaceUpload struct { envelope proto.Envelope request proto.WorkspaceWritePayload @@ -32,7 +32,8 @@ func (r *Router) handleWorkspaceWrite(ctx context.Context, env proto.Envelope) e // A malformed frame on an already admitted operation cannot claim that // its earlier commit did not execute. r.mu.Lock() - pending := r.workspaceWrite != nil && r.workspaceWrite.envelope.ID == env.ID + u := r.workspaceWrites[env.Assignment.SessionID] + pending := u != nil && u.envelope.ID == env.ID r.mu.Unlock() if pending { return errors.New("dispatch: malformed pending write frame") @@ -44,14 +45,14 @@ func (r *Router) handleWorkspaceWrite(ctx context.Context, env proto.Envelope) e r.mu.Unlock() return ErrRouterClosed } + u := r.workspaceWrites[env.Assignment.SessionID] if request.Step == "begin" { - if r.workspaceExport != nil && r.workspaceExport.request.Assignment.SessionID == request.SessionID || - r.runtimePreparation != nil && r.runtimePreparation.envelope.Assignment.SessionID == request.SessionID { + if r.workspaceExports[request.SessionID] != nil || r.runtimePreparations[request.SessionID] != nil { r.mu.Unlock() return r.sendWorkspaceWrite(ctx, env, rejectedWorkspaceWrite("resource_unavailable")) } - if r.workspaceWrite != nil { - duplicate := r.workspaceWrite.envelope.ID == env.ID + if u != nil || r.transferBytes > transferMemory-request.SizeBytes { + duplicate := u != nil && u.envelope.ID == env.ID r.mu.Unlock() if duplicate { return errors.New("dispatch: workspace write already admitted") @@ -73,14 +74,14 @@ func (r *Router) handleWorkspaceWrite(ctx context.Context, env proto.Envelope) e return r.sendWorkspaceWrite(ctx, env, rejectedWorkspaceWrite("resource_unavailable")) } u := &workspaceUpload{envelope: env, request: request, data: make([]byte, 0, request.SizeBytes), ready: make(chan struct{})} - r.workspaceWrite = u + r.workspaceWrites[request.SessionID] = u + r.transferBytes += request.SizeBytes done := r.trackWorkLocked(env.Assignment) r.shutdownWG.Add(1) r.mu.Unlock() go r.runWorkspaceUpload(context.WithoutCancel(ctx), u, environment, done) return r.sendWorkspaceWrite(ctx, env, proto.WorkspaceWriteResultPayload{Outcome: "ready"}) } - u := r.workspaceWrite if u == nil || u.envelope.ID != env.ID || u.envelope.Assignment != env.Assignment { r.mu.Unlock() return r.sendWorkspaceWrite(ctx, env, rejectedWorkspaceWrite("resource_unavailable")) @@ -132,9 +133,10 @@ func (r *Router) runWorkspaceUpload(ctx context.Context, u *workspaceUpload, env result = workspaceWriteResult(write, err, u.request.SizeBytes) } r.mu.Lock() + r.transferBytes -= u.request.SizeBytes u.uncertain = result.Outcome == "unknown" if !u.uncertain { - r.workspaceWrite = nil + delete(r.workspaceWrites, u.request.SessionID) } r.mu.Unlock() _ = r.sendWorkspaceWrite(ctx, u.envelope, result) diff --git a/docs/runtime-protocol.md b/docs/runtime-protocol.md index 857c81f81..7ec0d1897 100644 --- a/docs/runtime-protocol.md +++ b/docs/runtime-protocol.md @@ -110,17 +110,17 @@ An assignment binds one Session to the Runtime that runs it. `Envelope.assignmen Every Session frame carries the assignment: `execution_prepare`, `execution_start` and `execution_release`; `prompt_cancel`, `prompt_steer` and `function_result`; every frame of `runtime_prepare`, `workspace_read`, `workspace_write` and `workspace_export`; and `environment_quiesce` and `environment_resume`. A reply echoes its request's assignment, and a Run's frames carry the assignment that started it; Core rejects a reply or Run frame that names another. Heartbeats carry none. -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`. +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 and epoch with the same Environment, resource and grant is `bound` again; one with anything else fails with `assignment_conflict`. A bind of the bound assignment at a higher epoch supersedes the earlier epoch, with any Environment, resource or grant: the Runtime fences and cleans up the earlier epoch's work as a release does, keeping the home, then binds the new epoch and replies `bound`. Unfinished cleanup replies `failed` with `cleanup_unconfirmed`, and a retry at the same epoch repeats it; a release or a later bind in the meantime fails it with `assignment_stale`. 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`. -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`. An agent host's owner works in the Session's sandbox through File and Process on a [Link](./sandbox-link-protocol.md) attachment of its own, opened under the bind's grant on first use. It lasts from the Session's first bind until its home is removed and outlives the Session's Executors and connections. Quiescing the Runtime or releasing the assignment closes its attachment; a bind under a later assignment takes the owner over once that attachment is closed and otherwise fails with `assignment_conflict`. It runs each setup step as the Process operation whose ID is the `runtime_prepare` envelope ID. A `workspace_write` or `runtime_prepare` that cannot reach the sandbox before any effect ends `rejected` with `resource_unavailable`. It fails a Plugin whose MCP server declares literal `http_headers`, or is a stdio server with `env_vars`, before staging any of it. A file mutation or setup step whose outcome it cannot observe quarantines the owner until the home is removed: it sends no further mutation, and every later `workspace_write` and `runtime_prepare` ends `unknown`. 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`, `runtime_preparation_unsupported` for `runtime_prepare`, `write_unsupported` for `workspace_write`, and `read_unsupported` for `workspace_read` and `workspace_export`. +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`. An agent host's owner works in the Session's sandbox through File and Process on a [Link](./sandbox-link-protocol.md) attachment of its own, opened under the bind's grant on first use. It lasts from the Session's first bind until its home is removed and outlives the Session's Executors and connections. Quiescing its Environment, releasing the assignment or superseding its epoch closes its attachment; a bind under a later assignment or epoch takes the owner over once that attachment is closed, and a bind it cannot take over, including one of another Environment, fails with `assignment_conflict`. It runs each setup step as the Process operation whose ID is the `runtime_prepare` envelope ID. A `workspace_write` or `runtime_prepare` that cannot reach the sandbox before any effect ends `rejected` with `resource_unavailable`. It fails a Plugin whose MCP server declares literal `http_headers`, or is a stdio server with `env_vars`, before staging any of it. A file mutation or setup step whose outcome it cannot observe quarantines the owner until the home is removed: it sends no further mutation, and every later `workspace_write` and `runtime_prepare` ends `unknown`. 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`, `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. +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. `environment_quiesce` applies to the named Environment only: it fails with `resource_busy` while one of the Environment's Sessions has work in progress; otherwise the Runtime closes the Environment's Executors, has its owners release what they hold and replies `environment_quiesced`. Until the matching `environment_resume`, which carries the assignment that quiesced it, the Runtime admits for that Environment only the releases of its Sessions; Sessions of other Environments keep running. 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. ## Preparation and execution order -Environment initialization uses `runtime_prepare` on every connection, managed or user-owned; the [Environment contract](../contracts/agents-api/environments.md#runtime-capability-preparation) owns what is prepared and when. For a file or archive, send `begin`, wait for `ready`, send ordered chunks and await each matching `received` offset, then send `commit` and await `completed`. Initialization and finalization have typed headers without file data. Validate the expected outcome, offset, size and finite error code with the shared validator. One transfer is allowed per connection. A chunk receipt confirms staged bytes, not installation; a completed commit confirms that operation, not that a later Turn ran. +Environment initialization uses `runtime_prepare` on every connection, managed or user-owned; the [Environment contract](../contracts/agents-api/environments.md#runtime-capability-preparation) owns what is prepared and when. For a file or archive, send `begin`, wait for `ready`, send ordered chunks and await each matching `received` offset, then send `commit` and await `completed`. Initialization and finalization have typed headers without file data. Validate the expected outcome, offset, size and finite error code with the shared validator. A Session has at most one transfer of each kind at a time, and a transfer, or one whose outcome is unknown, blocks only its own Session's `execution_start` and transfers and its Environment's quiesce. The bodies staged by all of a connection's Sessions together are bounded by the sum of the `workspace_write` and `runtime_prepare` bounds; a `begin` beyond it is rejected with `write_capacity` or `runtime_preparation_capacity`. A chunk receipt confirms staged bytes, not installation; a completed commit confirms that operation, not that a later Turn ran. An execution Turn runs in five steps: @@ -202,7 +202,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. On every platform the owner lists directories, creates files and exports outputs in the daemon, 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 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 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 767f188c1..35c906bcf 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: 9d9f7867b0d5d88d0212ffaa9ef55ddf4aa5624a53717f2be92e0453f2049f01 +source_hash: a8ffdfb792698a153a1147afb1fb76fe19423085f05d6a1332e12bca278c8ed2 --- 此协议在 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)。 @@ -112,17 +112,17 @@ Usage frame 和最终 usage snapshot 都携带当前执行的累计测量,替 每个 Session frame 都携带分配:`execution_prepare`、`execution_start` 和 `execution_release`;`prompt_cancel`、`prompt_steer` 和 `function_result`;`runtime_prepare`、`workspace_read`、`workspace_write` 和 `workspace_export` 的每个 frame;以及 `environment_quiesce` 和 `environment_resume`。回复回显请求的分配,Run 的 frame 携带启动它的分配;Core 拒绝指明其他分配的回复或 Run frame。heartbeat 不携带分配。 -在一条连接上执行 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` 失败。 +在一条连接上执行 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` 失败。以相同的 epoch、Environment、resource 和 grant 重复绑定同一分配仍得到 `bound`;其他字段不同的同 epoch 绑定以 `assignment_conflict` 失败。以更高 epoch 绑定已绑定的分配会取代较早的 epoch,Environment、resource 和 grant 均可不同:Runtime 像释放那样 fence 并清理较早 epoch 的工作,但保留 home,然后绑定新 epoch 并回复 `bound`。清理未完成时回复 `failed` 和 `cleanup_unconfirmed`,以同一 epoch 重试会重复清理;期间到达的释放或更晚的绑定使其以 `assignment_stale` 失败。Runtime 只在其已绑定的分配下准入 Session frame:较旧的 epoch 或已释放的分配以 `assignment_stale` 失败;其他分配、Session 或 Environment 以 `assignment_conflict` 失败。已启动 Run 的 frame,包括其取消回执,在释放前仍可在启动它的分配下准入。Runtime 已记录回执的重复函数结果或决策只在应用它的分配下得到回答;其他分配以 `assignment_conflict` 失败。 -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 的绑定。agent host 的 owner 通过自己的一个 [Link](./sandbox-link-protocol.md) attachment,用 File 和 Process 在 Session 的沙箱中工作;该 attachment 在首次使用时凭绑定的 grant 打开。owner 从 Session 的第一次绑定存续到其 home 被删除,比 Session 的 Executor 和连接存续得更久。Runtime 静默(quiesce)或分配被释放时关闭其 attachment;之后分配下的绑定在该 attachment 关闭后接管 owner,否则以 `assignment_conflict` 失败。它把每个 setup 步骤作为 Process 操作运行,操作 ID 即 `runtime_prepare` 的 envelope ID。在产生任何作用前无法连到沙箱的 `workspace_write` 或 `runtime_prepare` 以 `rejected` 和 `resource_unavailable` 结束。若 Plugin 的 MCP server 声明了字面量 `http_headers`,或是带 `env_vars` 的 stdio server,owner 会在暂存其任何内容之前使其失败。无法观察到结果的文件变更或 setup 步骤会隔离 owner,直到 home 被删除:它不再发送任何变更,之后每个 `workspace_write` 和 `runtime_prepare` 都以 `unknown` 结束。没有 owner 的 Session 不支持上述任何操作,Runtime 以各自的类型化错误码拒绝:只读 preparation 为 `unsupported_read_preparation`;带 `local_environment` 的 Executor 配置为 `invalid_configuration`;`runtime_prepare` 为 `runtime_preparation_unsupported`;`workspace_write` 为 `write_unsupported`;`workspace_read` 和 `workspace_export` 为 `read_unsupported`。 +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 的绑定。agent host 的 owner 通过自己的一个 [Link](./sandbox-link-protocol.md) attachment,用 File 和 Process 在 Session 的沙箱中工作;该 attachment 在首次使用时凭绑定的 grant 打开。owner 从 Session 的第一次绑定存续到其 home 被删除,比 Session 的 Executor 和连接存续得更久。其 Environment 被 quiesce、分配被释放或其 epoch 被取代时关闭其 attachment;之后的分配或 epoch 下的绑定在该 attachment 关闭后接管 owner,它无法接管的绑定(包括其他 Environment 的绑定)以 `assignment_conflict` 失败。它把每个 setup 步骤作为 Process 操作运行,操作 ID 即 `runtime_prepare` 的 envelope ID。在产生任何作用前无法连到沙箱的 `workspace_write` 或 `runtime_prepare` 以 `rejected` 和 `resource_unavailable` 结束。若 Plugin 的 MCP server 声明了字面量 `http_headers`,或是带 `env_vars` 的 stdio server,owner 会在暂存其任何内容之前使其失败。无法观察到结果的文件变更或 setup 步骤会隔离 owner,直到 home 被删除:它不再发送任何变更,之后每个 `workspace_write` 和 `runtime_prepare` 都以 `unknown` 结束。没有 owner 的 Session 不支持上述任何操作,Runtime 以各自的类型化错误码拒绝:只读 preparation 为 `unsupported_read_preparation`;带 `local_environment` 的 Executor 配置为 `invalid_configuration`;`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 的分配。 +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 连接时重发所有未确认的释放;失败的释放退避重试,等待最久的释放先发送,因此失败的释放不会拖延其他释放。`environment_quiesce` 只作用于指定的 Environment:该 Environment 的某个 Session 仍有进行中的工作时,它以 `resource_busy` 失败;否则 Runtime 关闭该 Environment 的 Executor,让其 owner 释放所持有的资源,并回复 `environment_quiesced`。在匹配的 `environment_resume`(携带使其 quiesce 的分配)到达之前,Runtime 对该 Environment 只准入其 Session 的释放;其他 Environment 的 Session 继续运行。 Runtime 对无法路由的 Core frame 回复 `protocol_error`,回显请求 ID,并携带其类型和错误码。 ## 准备与执行顺序 {#preparation-and-execution-order} -无论托管还是用户自有环境,每条连接都通过 `runtime_prepare` 初始化 Environment;[Environment 契约](../../contracts/agents-api/zh/environments.md#runtime-capability-preparation)负责准备内容和时机。传输文件或 archive 时,发送 `begin`,等待 `ready`,发送有序 chunk 并等待每个匹配的 `received` offset,再发送 `commit` 并等待 `completed`。初始化和终结阶段使用不含文件数据的类型化 header。使用共享 validator 验证预期结果、offset、size 和有限错误 code。每条连接允许一个 transfer。chunk 回执确认暂存字节,不确认安装;完成的 commit 确认该操作,不证明后续 Turn 已运行。 +无论托管还是用户自有环境,每条连接都通过 `runtime_prepare` 初始化 Environment;[Environment 契约](../../contracts/agents-api/zh/environments.md#runtime-capability-preparation)负责准备内容和时机。传输文件或 archive 时,发送 `begin`,等待 `ready`,发送有序 chunk 并等待每个匹配的 `received` offset,再发送 `commit` 并等待 `completed`。初始化和终结阶段使用不含文件数据的类型化 header。使用共享 validator 验证预期结果、offset、size 和有限错误 code。每个 Session 同一时间每种 transfer 至多一个;一个 transfer,或结果未知的 transfer,只阻塞其所属 Session 的 `execution_start` 和 transfer,以及其 Environment 的 quiesce。一条连接上所有 Session 同时暂存的 body 总量以 `workspace_write` 与 `runtime_prepare` 限制之和为上限;超出上限的 `begin` 以 `write_capacity` 或 `runtime_preparation_capacity` 被拒绝。chunk 回执确认暂存字节,不确认安装;完成的 commit 确认该操作,不证明后续 Turn 已运行。 执行 Turn 分为五步: @@ -204,7 +204,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,直到精确结算;请求不重放。在各平台上,owner 都在 daemon 内列举目录、创建文件和导出输出,不使用外部 helper 或 staging directory。 +`workspace_write` 在原生 writer 运行前,通过已确认的 64 KiB frame 传输完整且有界的 body,验证声明的 digest,不运行模型。私有 transfer 限制为 50 MiB,与公开 API 在任何 Runtime 工作前检查的 5 MiB decoded inline 限制独立。Runtime 在接收或应用写入时排除该 Session 的执行;格式错误、不完整或到期的 transfer 不会到达 installer。精确的 commit 或拒绝回执释放 mutation owner。缺失或有歧义的回执保留不确定性:observer 取消和本地进程退出不能证明没有改变任何内容。公开准入前,Core 在 Session lock 下持久预约写入,跨重启阻止后继 mutation,直到精确结算;请求不重放。在各平台上,owner 都在 daemon 内列举目录、创建文件和导出输出,不使用外部 helper 或 staging directory。 ## MCP 连接权限 {#mcp-connection-authority} From 32be17da60062c41635d09614b2c601477a02fc3 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Wed, 7 Oct 2026 23:36:18 +0000 Subject: [PATCH 2/3] Serve Sessions from the agent-host subcommand --- .../agenthost/agenthost_linux_test.go | 9 +- apps/daemon/internal/agenthost/doc.go | 9 +- .../internal/agenthost/host_linux_test.go | 8 +- apps/daemon/internal/agenthost/host_other.go | 3 +- .../internal/agenthost/sessiondir_linux.go | 17 +- apps/daemon/internal/cli/agent_host_linux.go | 157 ++++++++++++++++++ .../internal/cli/agent_host_linux_test.go | 150 +++++++++++++++++ apps/daemon/internal/cli/agent_host_other.go | 14 ++ apps/daemon/internal/cli/connect.go | 57 ++++--- .../internal/cli/connect_cleanup_test.go | 3 +- apps/daemon/internal/cli/root.go | 1 + apps/daemon/internal/cli/root_test.go | 1 + docs/configuration.md | 9 +- docs/zh/configuration.md | 11 +- 14 files changed, 399 insertions(+), 50 deletions(-) create mode 100644 apps/daemon/internal/cli/agent_host_linux.go create mode 100644 apps/daemon/internal/cli/agent_host_linux_test.go create mode 100644 apps/daemon/internal/cli/agent_host_other.go diff --git a/apps/daemon/internal/agenthost/agenthost_linux_test.go b/apps/daemon/internal/agenthost/agenthost_linux_test.go index a1c92a984..05c3167e9 100644 --- a/apps/daemon/internal/agenthost/agenthost_linux_test.go +++ b/apps/daemon/internal/agenthost/agenthost_linux_test.go @@ -178,15 +178,8 @@ func newDaemon(t *testing.T, cfg Config, d deps) *daemon { func (dm *daemon) route(t *testing.T, reg *agent.Registry) { t.Helper() dm.frames, dm.opened = map[string]chan proto.Envelope{}, map[string]*session{} - removeHome := func(session string) error { - id, err := canonicalID(session) - if err != nil { - return err - } - return dm.host.RemoveHome(id) - } var err error - if dm.router, err = dispatch.New(dispatch.Config{Registry: reg, Sender: dm, Environments: dm.host.Environments, RemoveHome: removeHome, Log: slog.New(slog.DiscardHandler)}); err != nil { + if dm.router, err = dispatch.New(dispatch.Config{Registry: reg, Sender: dm, Environments: dm.host.Environments, RemoveHome: dm.host.RemoveHome, Log: slog.New(slog.DiscardHandler)}); err != nil { t.Fatal(err) } t.Cleanup(func() { dm.shutdown() }) diff --git a/apps/daemon/internal/agenthost/doc.go b/apps/daemon/internal/agenthost/doc.go index 90599726e..b409761cf 100644 --- a/apps/daemon/internal/agenthost/doc.go +++ b/apps/daemon/internal/agenthost/doc.go @@ -3,11 +3,10 @@ // Session's Link attachment. It needs Linux; elsewhere Open returns // ErrUnsupported. // -// The process that runs the agent host calls Open once at startup, hands -// Host.Registry and Host.Environments to the daemon's dispatch, which drives -// each Turn of each Executor, and calls Close once every Executor has closed. No production -// caller constructs Host.Registry yet; oac-daemon connect still runs -// Harnesses in the sandbox. Open takes two installation locks, which the +// The process that runs the agent host, oac-daemon agent-host, calls Open +// once at startup, hands Host.Registry and Host.Environments to the daemon's +// dispatch, which drives each Turn of each Executor, and calls Close once +// every Executor has closed. Open takes two installation locks, which the // Host holds until Close: a flock on StateDir/lock for the Session // directories and one on the Config.ViewCgroups directory for the view // cgroups. Under the locks, Open checks the requirements that need nothing diff --git a/apps/daemon/internal/agenthost/host_linux_test.go b/apps/daemon/internal/agenthost/host_linux_test.go index c7ef7ddf6..4df5d2f1a 100644 --- a/apps/daemon/internal/agenthost/host_linux_test.go +++ b/apps/daemon/internal/agenthost/host_linux_test.go @@ -16,6 +16,8 @@ import ( "syscall" "testing" + "github.com/google/uuid" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/sessionview" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/sessionview/sessionviewtest" @@ -193,7 +195,7 @@ func TestRemoveHome(t *testing.T) { opened <- e }() <-claimed - if err := h.RemoveHome(b.SessionID); !errors.Is(err, ErrSessionExists) { + if err := h.RemoveHome(uuid.UUID(b.SessionID).String()); !errors.Is(err, ErrSessionExists) { t.Fatalf("RemoveHome during an open = %v, want ErrSessionExists", err) } close(proceed) @@ -201,7 +203,7 @@ func TestRemoveHome(t *testing.T) { if e == nil { t.Fatal("open returned no Executor to close") } - if err := h.RemoveHome(b.SessionID); !errors.Is(err, ErrSessionExists) { + if err := h.RemoveHome(uuid.UUID(b.SessionID).String()); !errors.Is(err, ErrSessionExists) { t.Errorf("RemoveHome before the Executor closed = %v, want ErrSessionExists", err) } if _, err := os.Stat(f.session.Home.Host); err != nil { @@ -221,7 +223,7 @@ func TestRemoveHome(t *testing.T) { } // Removal is idempotent: an absent home is removed. for range 2 { - if err := h.RemoveHome(b.SessionID); err != nil { + if err := h.RemoveHome(uuid.UUID(b.SessionID).String()); err != nil { t.Fatalf("RemoveHome = %v", err) } } diff --git a/apps/daemon/internal/agenthost/host_other.go b/apps/daemon/internal/agenthost/host_other.go index 2a59377f8..63650e538 100644 --- a/apps/daemon/internal/agenthost/host_other.go +++ b/apps/daemon/internal/agenthost/host_other.go @@ -8,7 +8,6 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/dispatch" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" - "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" ) // Open reports that the agent host needs Linux. @@ -27,6 +26,6 @@ func (*Host) Environments(proto.AssignmentRef, proto.AssignmentBindPayload) disp } // RemoveHome reports that the agent host needs Linux. -func (*Host) RemoveHome(sandboxwire.ID) error { +func (*Host) RemoveHome(string) error { return fmt.Errorf("%w: remove home", ErrUnsupported) } diff --git a/apps/daemon/internal/agenthost/sessiondir_linux.go b/apps/daemon/internal/agenthost/sessiondir_linux.go index e827f8f4a..65d069c17 100644 --- a/apps/daemon/internal/agenthost/sessiondir_linux.go +++ b/apps/daemon/internal/agenthost/sessiondir_linux.go @@ -156,12 +156,17 @@ func sweep(stateDir string) error { return nil } -// RemoveHome drains and forgets the Session's Environment owner, then removes -// the Session's directory with its home once the Session's processes have -// settled. It returns ErrSessionExists while an Executor of the Session has -// not closed, and an Executor of the Session does not open while RemoveHome -// runs. A Session without a directory has nothing to remove. -func (h *Host) RemoveHome(id sandboxwire.ID) error { +// RemoveHome drains and forgets the Environment owner of the Session with +// the canonical ID session, then removes the Session's directory with its +// home once the Session's processes have settled, as +// dispatch.Config.RemoveHome. It returns ErrSessionExists while an Executor +// of the Session has not closed, and an Executor of the Session does not open +// while RemoveHome runs. A Session without a directory has nothing to remove. +func (h *Host) RemoveHome(session string) error { + id, err := canonicalID(session) + if err != nil { + return fmt.Errorf("%w: remove home: %w", ErrInvalidSession, err) + } if !claimSession(id) { return fmt.Errorf("%w: remove home", ErrSessionExists) } diff --git a/apps/daemon/internal/cli/agent_host_linux.go b/apps/daemon/internal/cli/agent_host_linux.go new file mode 100644 index 000000000..3abc06368 --- /dev/null +++ b/apps/daemon/internal/cli/agent_host_linux.go @@ -0,0 +1,157 @@ +//go:build linux + +package cli + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "log/slog" + "net/url" + "os" + "path/filepath" + "strings" + "time" + + "github.com/google/uuid" + "golang.org/x/sys/unix" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent/claudesdk" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent/codex" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent/mcode" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agenthost" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/daemonize" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/dispatch" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/paths" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/transport" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" +) + +// The agent-host image's layout (deploy/distribution/AgentHost.Dockerfile) +// and the container's own state. +const ( + agentHostManifest = "/opt/oac/harnesses.json" + agentHostShim = "/opt/oac/bin/oac-process-shim" + agentHostCADir = "/usr/share/ca-certificates/mozilla" + // agentHostState keeps each Session's home across restarts. + agentHostState = "/var/lib/oac/agent-host" + // agentHostCgroup is where the agent host mounts the container's own + // cgroup v2 hierarchy; its views directory is delegated to the views. + agentHostCgroup = "/run/oac/cgroup" + // agentHostUnreachable is how long Core may stay unreachable before the + // agent host exits, so that its supervisor restarts it. + agentHostUnreachable = 2 * time.Minute +) + +// agentHostUIDs is the range the agent host runs its Executors as. +var agentHostUIDs = agenthost.UIDRange{First: 70000, Count: 4096} + +func runAgentHost(rc *runContext, args []string) error { + return serveAgentHost(context.Background(), rc, args, harnessDeclarations) +} + +// serveAgentHost runs the agent host in its container: the Harnesses that +// the image's manifest installs and that declare a view serve Sessions bound +// to the Runtime of the identity, whose Environments are the agent host's. +func serveAgentHost(parent context.Context, rc *runContext, args []string, declarations []agent.Declaration) error { + flags := newFlagSet("agent-host") + identityFile := flags.String("identity-file", "", "path to the agent host's identity JSON") + coreURL := flags.String("core-url", "", "Core origin, https or a loopback http origin") + if err := flags.Parse(args); err != nil { + return fmt.Errorf("agent-host: parse flags: %w", err) + } + if flags.NArg() != 0 { + return fmt.Errorf("agent-host: unexpected arguments %q", flags.Args()) + } + var identity struct { + RuntimeID string `json:"runtime_id"` + Credential string `json:"credential"` + } + raw, err := os.ReadFile(*identityFile) + if err != nil { + return fmt.Errorf("agent-host: identity: %w", err) + } + decoder := json.NewDecoder(bytes.NewReader(raw)) + decoder.DisallowUnknownFields() + if err := decoder.Decode(&identity); err != nil { + return fmt.Errorf("agent-host: identity: %w", err) + } + runtimeID, err := uuid.Parse(identity.RuntimeID) + if err != nil || runtimeID.String() != identity.RuntimeID || identity.Credential == "" { + return errors.New("agent-host: identity needs a canonical runtime_id and a credential") + } + // Core derives the same URLs from OAC_PUBLIC_URL; its Link exists only + // on https or a loopback origin, which agenthost.Open checks. + origin, err := url.Parse(*coreURL) + if err != nil || (origin.Scheme != "http" && origin.Scheme != "https") || origin.Host == "" || origin.User != nil || + origin.Path != "" || origin.RawQuery != "" || origin.ForceQuery || origin.Fragment != "" || origin.String() != *coreURL { + return fmt.Errorf("agent-host: --core-url %q is not an http or https origin", *coreURL) + } + wsOrigin := "ws" + strings.TrimPrefix(*coreURL, "http") + + obslog.Init(obslog.Config{Format: "text", Level: slog.LevelInfo, Out: rc.stderr}) + ctx, cancel := daemonize.NotifyContext(parent) + defer cancel() + + env, err := agent.ManifestEnvironment(agentHostManifest, claudesdk.Installation(), codex.Installation(), mcode.Installation()) + if err != nil { + return fmt.Errorf("agent-host: %w", err) + } + for name, value := range env { + if err := os.Setenv(name, value); err != nil { + return fmt.Errorf("agent-host: %w", err) + } + } + harnesses := agent.NewRegistry() + for _, declaration := range declarations { + runtime := declaration.Discover(ctx, agent.DiscoveryOptions{Profile: paths.DefaultProfile, Stdout: rc.stdout, Stderr: rc.stderr}, declaration.Info) + if runtime == nil || runtime.View == nil { + continue + } + harnesses.RegisterKind(runtime.Info, declaration.Configuration) + harnesses.RegisterView(runtime.Info.Kind, *runtime.View) + } + + // With a private cgroup namespace, cgroup v2 mounted again is the + // container's own cgroup, writable unlike Docker's mount of it. + views := filepath.Join(agentHostCgroup, "views") + if err := os.MkdirAll(agentHostCgroup, 0o755); err != nil { + return fmt.Errorf("agent-host: %w", err) + } + if err := unix.Mount("cgroup2", agentHostCgroup, "cgroup2", 0, ""); err != nil { + return fmt.Errorf("%w: agent-host: mount cgroup v2: %w", agenthost.ErrUnsupported, err) + } + if err := os.Mkdir(views, 0o755); err != nil && !errors.Is(err, os.ErrExist) { + return fmt.Errorf("%w: agent-host: view cgroups: %w", agenthost.ErrUnsupported, err) + } + host, err := agenthost.Open(agenthost.Config{StateDir: agentHostState, UIDs: agentHostUIDs, ViewCgroups: views, + RelayURL: wsOrigin + "/api/v1/sandbox-link", RuntimeID: sandboxwire.ID(runtimeID), Credential: []byte(identity.Credential), + Harnesses: harnesses, Shim: agentHostShim, CADir: agentHostCADir, Log: obslog.Bg()}) + if err != nil { + return fmt.Errorf("agent-host: %w", err) + } + defer host.Close() + + // Bootstrap on each dial, so a Core that is still starting is retried + // with the connection's backoff. The agent host dials Core's origin, not + // the bootstrap's public ws_url. + wsURL := wsOrigin + "/api/v1/agent-daemon/ws" + var boot *transport.BootstrapResponse + dial := func(ctx context.Context) (*transport.Conn, error) { + b, err := transport.Bootstrap(ctx, *coreURL+"/api/v1", identity.RuntimeID, identity.Credential, Version) + if err != nil { + return nil, err + } + boot = b + return transport.Dial(ctx, transport.DialOptions{WSURL: wsURL, DeviceID: identity.RuntimeID, Credential: identity.Credential, DaemonVersion: proto.Version}) + } + cfg := dispatch.Config{Registry: host.Registry(), Environments: host.Environments, RemoveHome: host.RemoveHome} + return serveConnections(ctx, wsURL, dial, agentHostUnreachable, func(conn *transport.Conn) error { + return pumpConn(ctx, conn, cfg, boot) + }) +} diff --git a/apps/daemon/internal/cli/agent_host_linux_test.go b/apps/daemon/internal/cli/agent_host_linux_test.go new file mode 100644 index 000000000..32fd13362 --- /dev/null +++ b/apps/daemon/internal/cli/agent_host_linux_test.go @@ -0,0 +1,150 @@ +//go:build linux + +package cli + +import ( + "context" + "encoding/json" + "encoding/pem" + "errors" + "io" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "slices" + "strings" + "testing" + "time" + + "github.com/google/uuid" + "github.com/gorilla/websocket" + "golang.org/x/sys/unix" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/transport" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto/prototest" +) + +// TestAgentHostReportsItsDeclarations runs the agent host as its container +// does, against a Core peer, with an image that installs no Harness: once +// without declarations and once with two, of which one has a view. The first +// heartbeat declares exactly the kinds with a view, and home removal. +func TestAgentHostReportsItsDeclarations(t *testing.T) { + if os.Getenv("OAC_TEST_AGENTHOST") != "1" { + t.Skip("set OAC_TEST_AGENTHOST=1 and run the test as root in a throwaway container with the agent-host container's flags") + } + issuer := httptest.NewTLSServer(nil) + issuer.Close() + for path, content := range map[string][]byte{ + agentHostManifest: []byte(`{"node": "/usr/local/bin/node", "harnesses": {}}`), + filepath.Join(agentHostCADir, "ca.crt"): pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: issuer.Certificate().Raw}), + } { + if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(path, content, 0o644); err != nil { + t.Fatal(err) + } + } + // The agent host presents the credential as written; it is not decoded. + runtimeID, credential := uuid.NewString(), "c2VjcmV0K/8=" + identity := filepath.Join(t.TempDir(), "identity.json") + if err := os.WriteFile(identity, []byte(`{"runtime_id": "`+runtimeID+`", "credential": "`+credential+`"}`), 0o600); err != nil { + t.Fatal(err) + } + + heartbeats := make(chan proto.HeartbeatPayload, 1) + core := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Header.Get("Authorization") != "Bearer "+credential { + http.Error(w, "credential", http.StatusUnauthorized) + return + } + switch r.URL.Path { + case "/api/v1/agent-daemon/bootstrap": + _ = json.NewEncoder(w).Encode(transport.BootstrapResponse{DeviceID: runtimeID, HeartbeatSeconds: 60}) + case "/api/v1/agent-daemon/ws": + if r.URL.Query().Get("device_id") != runtimeID { + http.Error(w, "device", http.StatusUnauthorized) + return + } + peer, err := (&websocket.Upgrader{}).Upgrade(w, r, nil) + if err != nil { + return + } + defer peer.Close() + for { + var env proto.Envelope + if err := peer.ReadJSON(&env); err != nil { + return + } + var heartbeat proto.HeartbeatPayload + if env.Type == proto.TypeHeartbeat && env.DecodePayload(&heartbeat) == nil { + select { + case heartbeats <- heartbeat: + default: + } + } + } + default: + http.NotFound(w, r) + } + })) + defer core.Close() + + declare := func(kind string, view *agent.View) agent.Declaration { + info := proto.SupportedAgentKind{Kind: kind, Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{ + LocalEnvironment: proto.CapabilitySupported, EnvironmentNone: proto.CapabilitySupported})} + return agent.Declaration{Info: info, Configuration: prototest.ModelConfiguration(), + Discover: func(context.Context, agent.DiscoveryOptions, proto.SupportedAgentKind) *agent.Runtime { + return &agent.Runtime{Info: info, View: view} + }} + } + view := &agent.View{Proxy: agent.ViewProxyEnv, Executor: func(context.Context, proto.PromptRequestPayload, agent.ViewSession) (agent.Executor, error) { + return nil, errors.New("no Executor") + }} + // The image may install no Harness: the agent host still connects and + // declares no kind. + for _, c := range []struct { + declarations []agent.Declaration + kinds []string + }{{nil, nil}, {[]agent.Declaration{declare("viewed", view), declare("unviewed", nil)}, []string{"viewed"}}} { + ctx, cancel := context.WithCancel(t.Context()) + served := make(chan error, 1) + go func() { + rc := &runContext{stdin: strings.NewReader(""), stdout: io.Discard, stderr: os.Stderr} + served <- serveAgentHost(ctx, rc, []string{"--identity-file", identity, "--core-url", core.URL}, c.declarations) + }() + select { + case heartbeat := <-heartbeats: + var kinds []string + for _, kind := range heartbeat.SupportedAgentKinds { + kinds = append(kinds, kind.Kind) + } + if !slices.Equal(kinds, c.kinds) || heartbeat.HomeRemoval != proto.CapabilitySupported { + t.Errorf("heartbeat declares %q with home removal %q, want %q with home removal", kinds, heartbeat.HomeRemoval, c.kinds) + } + case err := <-served: + t.Fatalf("agent host stopped before its first heartbeat: %v", err) + case <-time.After(30 * time.Second): + t.Fatal("no heartbeat") + } + cancel() + if err := <-served; err != nil { + t.Fatalf("agent host stopped with %v, want nil after its signal", err) + } + if err := unix.Unmount(agentHostCgroup, 0); err != nil { + t.Fatal(err) + } + } +} + +// TestAgentHostExitsWhileCoreStaysUnreachable checks the bound after which +// the agent host exits for its supervisor to restart it. +func TestAgentHostExitsWhileCoreStaysUnreachable(t *testing.T) { + dial := func(context.Context) (*transport.Conn, error) { return nil, errors.New("connection refused") } + if err := serveConnections(t.Context(), "ws://127.0.0.1:1", dial, 10*time.Millisecond, nil); !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("serveConnections = %v, want the unreachable bound", err) + } +} diff --git a/apps/daemon/internal/cli/agent_host_other.go b/apps/daemon/internal/cli/agent_host_other.go new file mode 100644 index 000000000..0b75b8eb4 --- /dev/null +++ b/apps/daemon/internal/cli/agent_host_other.go @@ -0,0 +1,14 @@ +//go:build !linux + +package cli + +import ( + "fmt" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agenthost" +) + +// runAgentHost reports that the agent host needs Linux. +func runAgentHost(*runContext, []string) error { + return fmt.Errorf("%w: agent-host", agenthost.ErrUnsupported) +} diff --git a/apps/daemon/internal/cli/connect.go b/apps/daemon/internal/cli/connect.go index 157566dde..70bb9ef1b 100644 --- a/apps/daemon/internal/cli/connect.go +++ b/apps/daemon/internal/cli/connect.go @@ -263,12 +263,25 @@ func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof if control != nil { return runSuspendLoop(rootCtx, dial, registry, local, boot, agentCLIs, control) } + return serveConnections(rootCtx, wsURL, dial, 0, func(conn *transport.Conn) error { + return pumpConn(rootCtx, conn, dispatch.Config{Registry: registry, Environments: localEnvironments(local)}, boot) + }) +} + +// serveConnections dials Core with dial and serves each connection until ctx +// ends or Core rejects the credential. With a positive unreachable bound it +// fails once Core has stayed unreachable that long. +func serveConnections(ctx context.Context, wsURL string, dial transport.DialFn, unreachable time.Duration, serve func(*transport.Conn) error) error { for { - if err := rootCtx.Err(); err != nil { + if err := ctx.Err(); err != nil { return nil } - conn, err := transport.Reconnect(rootCtx, dial, transport.DefaultBackoff, func(attempt int, lastDelay time.Duration, lastErr error) { + dialCtx, stop := ctx, context.CancelFunc(func() {}) + if unreachable > 0 { + dialCtx, stop = context.WithTimeout(ctx, unreachable) + } + conn, err := transport.Reconnect(dialCtx, dial, transport.DefaultBackoff, func(attempt int, lastDelay time.Duration, lastErr error) { switch { case attempt == 1: obslog.Bg().Info("connecting", "ws_url", wsURL) @@ -281,21 +294,24 @@ func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof obslog.Bg().Warn("dial retry", "attempt", attempt, "delay", lastDelay) } }) + stop() if err != nil { - if errors.Is(err, context.Canceled) { + if ctx.Err() != nil { return nil } if errors.Is(err, transport.ErrPermanent) { return fmt.Errorf("connect: permanent error (reissue the daemon credential): %w", err) } + if errors.Is(err, context.DeadlineExceeded) { + return fmt.Errorf("connect: Core unreachable for %s: %w", unreachable, err) + } return fmt.Errorf("connect: dial: %w", err) } obslog.Bg().Info("ws connected", "device_id", conn.DeviceID()) - // pumpConn returns on conn close (peer hangup, transport - // error, root ctx cancel). Loop back into Reconnect unless - // root ctx is cancelled. - pumpErr := pumpConn(rootCtx, conn, registry, local, boot, agentCLIs) + // serve returns on conn close (peer hangup, transport error, ctx + // cancel). Loop back into Reconnect unless ctx is cancelled. + pumpErr := serve(conn) if pumpErr != nil { obslog.Bg().Warn("ws session ended", "err", pumpErr) } else { @@ -305,7 +321,7 @@ func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof // Server-initiated clean close (e.g. shutdown) → exit; // otherwise loop back and reconnect. - if rootCtx.Err() != nil { + if ctx.Err() != nil { return nil } // Permanent error (e.g. runtime deleted) → exit instead of @@ -315,7 +331,7 @@ func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof } // Small breather before redialing so a flapping server doesn't // get a tight loop of upgrade requests. - _ = transport.Sleep(rootCtx, 1*time.Second) + _ = transport.Sleep(ctx, 1*time.Second) } } @@ -328,17 +344,14 @@ func localEnvironments(local *localworkspace.Binding) func(proto.AssignmentRef, return local.Resolve } -// 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. -// Failed cleanup keeps this exact Router alive, including after a shutdown signal. -func pumpConn(parentCtx context.Context, conn *transport.Conn, registry *agent.Registry, local *localworkspace.Binding, boot *transport.BootstrapResponse, agentCLIs agentCLIDiscovery) error { - router, err := dispatch.New(dispatch.Config{ - Registry: registry, - Sender: conn, - Log: obslog.Bg(), - Environments: localEnvironments(local), - }) +// pumpConn runs the per-connection workload: a dispatch.Router of cfg's +// Harness kinds and Environment owners fed by conn.Recv(), heartbeats every +// boot.HeartbeatInterval(), and a confirmed router.Shutdown before returning +// ownership to the reconnect loop. Failed cleanup keeps this exact Router +// alive, including after a shutdown signal. +func pumpConn(parentCtx context.Context, conn *transport.Conn, cfg dispatch.Config, boot *transport.BootstrapResponse) error { + cfg.Sender, cfg.Log = conn, obslog.Bg() + router, err := dispatch.New(cfg) if err != nil { return fmt.Errorf("router init: %w", err) } @@ -352,8 +365,8 @@ func pumpConn(parentCtx context.Context, conn *transport.Conn, registry *agent.R Timestamp: time.Now().Unix(), ActiveRequests: router.ActiveRuns(), DaemonVersion: Version, - SupportedAgentKinds: registry.SupportedAgentKinds(), - HomeRemoval: proto.CapabilityUnsupported, + SupportedAgentKinds: cfg.Registry.SupportedAgentKinds(), + HomeRemoval: proto.CapabilityFromBool(cfg.RemoveHome != nil), } }, obslog.Bg().With("component", "heartbeat")) diff --git a/apps/daemon/internal/cli/connect_cleanup_test.go b/apps/daemon/internal/cli/connect_cleanup_test.go index d707804f6..bd1c1d4e7 100644 --- a/apps/daemon/internal/cli/connect_cleanup_test.go +++ b/apps/daemon/internal/cli/connect_cleanup_test.go @@ -9,6 +9,7 @@ import ( "strings" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/dispatch" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/transport" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto/prototest" @@ -133,7 +134,7 @@ func testDisconnectedPumpCleanup(t *testing.T, suspend bool) { return } defer conn.Close() - finished <- pumpConn(ctx, conn, registry, nil, boot, agentCLIDiscovery{}) + finished <- pumpConn(ctx, conn, dispatch.Config{Registry: registry}, boot) }() var peer *websocket.Conn select { diff --git a/apps/daemon/internal/cli/root.go b/apps/daemon/internal/cli/root.go index cd77a7d2b..2b3323e33 100644 --- a/apps/daemon/internal/cli/root.go +++ b/apps/daemon/internal/cli/root.go @@ -41,6 +41,7 @@ var commands = []command{ {name: "runtime-mcp-exec", summary: "Execute an installed MCP server", run: runRuntimeMCP}, {name: "placement", summary: "Enroll or retire an explicitly managed local execution placement", run: runPlacement}, {name: "connect", summary: "Open the reverse WebSocket and start serving prompts", run: runConnect}, + {name: "agent-host", summary: "Serve Sessions from the agent-host container", run: runAgentHost}, {name: "status", summary: "Print the credential profile and daemon state", run: runStatus}, {name: "stop", summary: "Stop a background `connect -b` daemon", run: runStop}, {name: "logs", summary: "Tail the background daemon's log file", run: runLogs}, diff --git a/apps/daemon/internal/cli/root_test.go b/apps/daemon/internal/cli/root_test.go index c27176d35..b235f3d2f 100644 --- a/apps/daemon/internal/cli/root_test.go +++ b/apps/daemon/internal/cli/root_test.go @@ -66,6 +66,7 @@ func TestSubcommandsAreRegistered(t *testing.T) { "runtime-mcp-exec": false, "placement": false, "connect": false, + "agent-host": false, "status": false, "stop": false, "logs": false, diff --git a/docs/configuration.md b/docs/configuration.md index ae09a0fba..ded27d690 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -146,7 +146,14 @@ The agent host runs each Session's Harness outside the sandbox, in a view of its Docker's default seccomp profile stays: with `CAP_SYS_ADMIN` it allows `clone3`, `mount` and `unshare`. The container gets no Docker socket and publishes no port. -The agent host needs a cgroup v2 directory delegated to it. It starts each view in a cgroup of its own there, ends the view with `cgroup.kill` and removes the cgroup. When it starts it ends and removes every cgroup in the directory, because each counts as a view's, so nothing else may use it. The directory must be writable and must not contain the agent host's own process. Docker mounts the container's cgroup read-only. With the flags above, cgroup v2 mounted again inside the container (`mount -t cgroup2 cgroup2